Merge branch 'bool-api'

This commit is contained in:
weisd
2024-11-19 10:37:31 +08:00
16 changed files with 550 additions and 161 deletions
Generated
+1
View File
@@ -2115,6 +2115,7 @@ dependencies = [
"s3s",
"serde",
"serde-xml-rs",
"serde_json",
"serde_urlencoded",
"time",
"tracing",
+1 -1
View File
@@ -14,7 +14,7 @@ use super::condition::{
#[derive(Debug, Deserialize, Serialize, Default, Clone, PartialEq, Eq)]
pub struct ActionSet(HashSet<Action>);
pub struct ActionSet(pub HashSet<Action>);
impl ActionSet {
pub fn is_match(&self, act: &Action) -> bool {
+4 -3
View File
@@ -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,
}
@@ -106,6 +106,9 @@ impl Functions {
}
set
}
pub fn is_empty(&self) -> bool {
self.0.is_empty()
}
}
impl Debug for Functions {
+11 -1
View File
@@ -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<Resource>);
pub struct ResourceSet(pub HashSet<Resource>);
impl ResourceSet {
pub fn validate_bucket(&self, bucket: &str) -> Result<()> {
@@ -222,6 +228,10 @@ impl ResourceSet {
}
false
}
pub fn is_empty(&self) -> bool {
self.0.is_empty()
}
}
impl AsRef<HashSet<Resource>> for ResourceSet {
+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));
}
@@ -766,9 +766,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)
}
},
}
}
@@ -809,13 +827,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 @@ pub 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!()
}
+48 -3
View File
@@ -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,
}
@@ -89,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() {
@@ -152,6 +154,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 +289,48 @@ 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 {
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());
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 {
disks,
..Default::default()
}
}
async fn list_bucket(&self, _opts: &BucketOptions) -> Result<Vec<BucketInfo>> {
unimplemented!()
}
+75 -5
View File
@@ -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,8 @@ 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::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 +54,7 @@ use std::{
};
use time::OffsetDateTime;
use tokio::fs;
use tokio::sync::Semaphore;
use tokio::sync::{RwLock, Semaphore};
use tracing::{debug, info, warn};
use uuid::Uuid;
@@ -67,6 +69,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 +170,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;
@@ -648,9 +654,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());
}
}
@@ -814,6 +818,72 @@ lazy_static! {
#[async_trait::async_trait]
impl StorageAPI for ECStore {
async fn backend_info(&self) -> BackendInfo {
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)
.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_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_sc_parity {
standard_sc_data.push(set_count - 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);
}
BackendInfo {
backend_type: BackendByte::Erasure,
online_disks: BackendDisks::new(),
offline_disks: BackendDisks::new(),
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 {
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: backend, disks }
}
async fn list_bucket(&self, opts: &BucketOptions) -> Result<Vec<BucketInfo>> {
// TODO: opts.cached
+82 -3
View File
@@ -780,6 +780,84 @@ pub struct DeletedObject {
// pub replication_state: ReplicationState,
}
#[derive(Debug, Default, Serialize)]
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<String>,
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<DiskMetrics>,
// pub heal_info: Option<HealingDisk>,
pub used_inodes: u64,
pub free_inodes: u64,
pub local: bool,
pub pool_index: i32,
pub set_index: i32,
pub disk_index: i32,
}
#[derive(Debug, Default)]
pub struct StorageInfo {
pub disks: Vec<StorageDisk>,
pub backend: BackendInfo,
}
#[derive(Debug, Default, Serialize)]
pub struct BackendDisks(HashMap<String, usize>);
impl BackendDisks {
pub fn new() -> Self {
Self(HashMap::new())
}
pub fn sum(&self) -> usize {
self.0.values().sum()
}
}
#[derive(Debug, Default, Serialize)]
#[serde(rename_all = "PascalCase", default)]
pub struct BackendInfo {
pub backend_type: BackendByte,
pub online_disks: BackendDisks,
pub offline_disks: BackendDisks,
#[serde(rename = "StandardSCData")]
pub standard_sc_data: Vec<usize>,
#[serde(rename = "StandardSCParities")]
pub standard_sc_parities: Vec<usize>,
#[serde(rename = "StandardSCParity")]
pub standard_sc_parity: Option<usize>,
#[serde(rename = "RRSCData")]
pub rr_sc_data: Vec<usize>,
#[serde(rename = "RRSCParities")]
pub rr_sc_parities: Vec<usize>,
#[serde(rename = "RRSCParity")]
pub rr_sc_parity: Option<usize>,
pub total_sets: Vec<usize>,
pub drives_per_set: Vec<usize>,
}
#[async_trait::async_trait]
pub trait ObjectIO: Send + Sync + 'static {
// GetObjectNInfo
@@ -800,9 +878,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<BucketInfo>;
+1
View File
@@ -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
+62 -4
View File
@@ -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.clone()),
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))))
}
}