Merge pull request #132 from rustfs/peer-rest

Peer rest
This commit is contained in:
junxiangMu
2024-11-27 15:01:14 +08:00
committed by GitHub
28 changed files with 4694 additions and 57 deletions
Generated
+98 -1
View File
@@ -541,6 +541,26 @@ dependencies = [
"rand_core",
]
[[package]]
name = "darwin-libproc"
version = "0.1.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "9fb90051930c9a0f09e585762152048e23ac74d20c10590ef7cf01c0343c3046"
dependencies = [
"darwin-libproc-sys",
"libc",
"memchr",
]
[[package]]
name = "darwin-libproc-sys"
version = "0.1.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "57cebb5bde66eecdd30ddc4b9cd208238b15db4982ccc72db59d699ea10867c1"
dependencies = [
"libc",
]
[[package]]
name = "deranged"
version = "0.3.11"
@@ -551,6 +571,17 @@ dependencies = [
"serde",
]
[[package]]
name = "derive_more"
version = "0.99.18"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "5f33878137e4dafd7fa914ad4e259e18a4e8e532b9617a2d0150262bf53abfce"
dependencies = [
"proc-macro2",
"quote",
"syn",
]
[[package]]
name = "digest"
version = "0.10.7"
@@ -594,6 +625,8 @@ dependencies = [
"lazy_static",
"lock",
"protos",
"rmp-serde",
"serde",
"serde_json",
"tokio",
"tonic",
@@ -614,6 +647,7 @@ dependencies = [
"bytesize",
"common",
"crc32fast",
"flatbuffers",
"futures",
"glob",
"hex-simd",
@@ -623,7 +657,7 @@ dependencies = [
"lock",
"md-5",
"netif",
"nix",
"nix 0.29.0",
"num",
"num_cpus",
"openssl",
@@ -1407,6 +1441,24 @@ dependencies = [
"hashbrown 0.12.3",
]
[[package]]
name = "mach2"
version = "0.4.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "19b955cdeb2a02b9117f121ce63aa52d08ade45de53e48fe6a38b39c10f6f709"
dependencies = [
"libc",
]
[[package]]
name = "madmin"
version = "0.0.1"
dependencies = [
"ecstore",
"psutil",
"serde",
]
[[package]]
name = "matchers"
version = "0.1.0"
@@ -1493,6 +1545,17 @@ dependencies = [
"winapi",
]
[[package]]
name = "nix"
version = "0.24.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "fa52e972a9a719cecb6864fb88568781eb706bac2cd1d4f04a648542dbf78069"
dependencies = [
"bitflags 1.3.2",
"cfg-if",
"libc",
]
[[package]]
name = "nix"
version = "0.29.0"
@@ -1812,6 +1875,12 @@ version = "0.3.31"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "953ec861398dccce10c670dfeaf3ec4911ca479e9c02154b3a215178c5f566f2"
[[package]]
name = "platforms"
version = "2.0.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "e8d0eef3571242013a0d5dc84861c3ae4a652e56e12adf8bdc26ff5f8cb34c94"
[[package]]
name = "powerfmt"
version = "0.2.0"
@@ -1934,6 +2003,25 @@ dependencies = [
"tower 0.5.1",
]
[[package]]
name = "psutil"
version = "3.3.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "5e617cc9058daa5e1fe5a0d23ed745773a5ee354111dad1ec0235b0cc16b6730"
dependencies = [
"cfg-if",
"darwin-libproc",
"derive_more",
"glob",
"mach2",
"nix 0.24.3",
"num_cpus",
"once_cell",
"platforms",
"thiserror",
"unescape",
]
[[package]]
name = "quick-xml"
version = "0.37.0"
@@ -2143,6 +2231,7 @@ dependencies = [
"lazy_static",
"lock",
"log",
"madmin",
"matchit 0.8.5",
"mime",
"netif",
@@ -2152,6 +2241,7 @@ dependencies = [
"prost-types",
"protobuf",
"protos",
"rmp-serde",
"s3s",
"serde",
"serde_json",
@@ -2169,6 +2259,7 @@ dependencies = [
"tracing-error",
"tracing-subscriber",
"transform-stream",
"url",
"uuid",
]
@@ -2896,6 +2987,12 @@ dependencies = [
"tz-rs",
]
[[package]]
name = "unescape"
version = "0.1.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "ccb97dac3243214f8d8507998906ca3e2e0b900bf9bf4870477f125b82e68f6e"
[[package]]
name = "unicode-ident"
version = "1.0.13"
+4
View File
@@ -1,6 +1,7 @@
[workspace]
resolver = "2"
members = [
"madmin",
"rustfs",
"ecstore",
"e2e_test",
@@ -20,6 +21,7 @@ rust-version = "1.75"
version = "0.0.1"
[workspace.dependencies]
madmin = { path = "./madmin" }
async-trait = "0.1.83"
backon = "1.2.0"
bytes = "1.8.0"
@@ -51,6 +53,8 @@ prost-types = "0.13.3"
protobuf = "3.7"
protos = { path = "./common/protos" }
rand = "0.8.5"
rmp = "0.8.14"
rmp-serde = "1.3.0"
s3s = { git = "https://github.com/Nugine/s3s.git", rev = "207170f526c75a8190e8f9afadf961909fd01d34", default-features = true, features = [
"tower",
] }
File diff suppressed because it is too large Load Diff
+360 -2
View File
@@ -411,8 +411,325 @@ message GenerallyLockRequest {
}
message GenerallyLockResponse {
bool success = 1;
optional string error_info = 2;
bool success = 1;
optional string error_info = 2;
}
message Mss {
map<string, string> value = 1;
}
message LocalStorageInfoRequest {
bool metrics = 1;
}
message LocalStorageInfoResponse {
bool success = 1;
bytes storage_info = 2;
optional string error_info = 3;
}
message ServerInfoRequest {
bool metrics = 1;
}
message ServerInfoResponse {
bool success = 1;
bytes server_properties = 2;
optional string error_info = 3;
}
message GetCpusRequest {}
message GetCpusResponse {
bool success = 1;
bytes cpus = 2;
optional string error_info = 3;
}
message GetNetInfoRequest {}
message GetNetInfoResponse {
bool success = 1;
bytes net_info = 2;
optional string error_info = 3;
}
message GetPartitionsRequest {}
message GetPartitionsResponse {
bool success = 1;
bytes partitions = 2;
optional string error_info = 3;
}
message GetOsInfoRequest {}
message GetOsInfoResponse {
bool success = 1;
bytes os_info = 2;
optional string error_info = 3;
}
message GetSELinuxInfoRequest {}
message GetSELinuxInfoResponse {
bool success = 1;
bytes sys_services = 2;
optional string error_info = 3;
}
message GetSysConfigRequest {}
message GetSysConfigResponse {
bool success = 1;
bytes sys_config = 2;
optional string error_info = 3;
}
message GetSysErrorsRequest {}
message GetSysErrorsResponse {
bool success = 1;
bytes sys_errors = 2;
optional string error_info = 3;
}
message GetMemInfoRequest {}
message GetMemInfoResponse {
bool success = 1;
bytes mem_info = 2;
optional string error_info = 3;
}
message GetMetricsRequest {
uint64 metric_type = 1;
bytes opts = 2;
}
message GetMetricsResponse {
bool success = 1;
bytes realtime_metrics = 2;
optional string error_info = 3;
}
message GetProcInfoRequest {}
message GetProcInfoResponse {
bool success = 1;
bytes proc_info = 2;
optional string error_info = 3;
}
message StartProfilingRequest {
string profiler = 1;
}
message StartProfilingResponse {
bool success = 1;
optional string error_info = 2;
}
message DownloadProfileDataRequest {}
message DownloadProfileDataResponse {
bool success = 1;
map<string, bytes> data = 2;
optional string error_info = 3;
}
message GetBucketStatsDataRequest {
string bucket = 1;
}
message GetBucketStatsDataResponse {
bool success = 1;
bytes bucket_stats = 2;
optional string error_info = 3;
}
message GetSRMetricsDataRequest {}
message GetSRMetricsDataResponse {
bool success = 1;
bytes sr_metrics_summary = 2;
optional string error_info = 3;
}
message GetAllBucketStatsRequest {}
message GetAllBucketStatsResponse {
bool success = 1;
bytes bucket_stats_map = 2;
optional string error_info = 3;
}
message LoadBucketMetadataRequest {
string bucket = 1;
}
message LoadBucketMetadataResponse {
bool success = 1;
optional string error_info = 2;
}
message DeleteBucketMetadataRequest {
string bucket = 1;
}
message DeleteBucketMetadataResponse {
bool success = 1;
optional string error_info = 2;
}
message DeletePolicyRequest {
string policy_name = 1;
}
message DeletePolicyResponse {
bool success = 1;
optional string error_info = 2;
}
message LoadPolicyRequest {
string policy_name = 1;
}
message LoadPolicyResponse {
bool success = 1;
optional string error_info = 2;
}
message LoadPolicyMappingRequest {
string user_or_group = 1;
uint64 user_type = 2;
bool is_group = 3;
}
message LoadPolicyMappingResponse {
bool success = 1;
optional string error_info = 2;
}
message DeleteUserRequest {
string access_key = 1;
}
message DeleteUserResponse {
bool success = 1;
optional string error_info = 2;
}
message DeleteServiceAccountRequest {
string access_key = 1;
}
message DeleteServiceAccountResponse {
bool success = 1;
optional string error_info = 2;
}
message LoadUserRequest {
string access_key = 1;
bool temp = 2;
}
message LoadUserResponse {
bool success = 1;
optional string error_info = 2;
}
message LoadServiceAccountRequest {
string access_key = 1;
}
message LoadServiceAccountResponse {
bool success = 1;
optional string error_info = 2;
}
message LoadGroupRequest {
string group = 1;
}
message LoadGroupResponse {
bool success = 1;
optional string error_info = 2;
}
message ReloadSiteReplicationConfigRequest {}
message ReloadSiteReplicationConfigResponse {
bool success = 1;
optional string error_info = 2;
}
message SignalServiceRequest {
Mss vars = 1;
}
message SignalServiceResponse {
bool success = 1;
optional string error_info = 2;
}
message BackgroundHealStatusRequest {}
message BackgroundHealStatusResponse {
bool success = 1;
bytes bg_heal_state = 2;
optional string error_info = 3;
}
message GetMetacacheListingRequest {
bytes opts = 1;
}
message GetMetacacheListingResponse {
bool success = 1;
bytes metacache = 2;
optional string error_info = 3;
}
message UpdateMetacacheListingRequest {
bytes metacache = 1;
}
message UpdateMetacacheListingResponse {
bool success = 1;
bytes metacache = 2;
optional string error_info = 3;
}
message ReloadPoolMetaRequest {}
message ReloadPoolMetaResponse {
bool success = 1;
optional string error_info = 2;
}
message StopRebalanceRequest {}
message StopRebalanceResponse {
bool success = 1;
optional string error_info = 2;
}
message LoadRebalanceMetaRequest {
bool start_rebalance = 1;
}
message LoadRebalanceMetaResponse {
bool success = 1;
optional string error_info = 2;
}
message LoadTransitionTierConfigRequest {}
message LoadTransitionTierConfigResponse {
bool success = 1;
optional string error_info = 2;
}
/* -------------------------------------------------------------------- */
@@ -466,4 +783,45 @@ service NodeService {
rpc RUnLock(GenerallyLockRequest) returns (GenerallyLockResponse) {};
rpc ForceUnLock(GenerallyLockRequest) returns (GenerallyLockResponse) {};
rpc Refresh(GenerallyLockRequest) returns (GenerallyLockResponse) {};
/* -------------------------------peer rest service-------------------------- */
rpc LocalStorageInfo(LocalStorageInfoRequest) returns (LocalStorageInfoResponse) {};
rpc ServerInfo(ServerInfoRequest) returns (ServerInfoResponse) {};
rpc GetCpus(GetCpusRequest) returns (GetCpusResponse) {};
rpc GetNetInfo(GetNetInfoRequest) returns (GetNetInfoResponse) {};
rpc GetPartitions(GetPartitionsRequest) returns (GetPartitionsResponse) {};
rpc GetOsInfo(GetOsInfoRequest) returns (GetOsInfoResponse) {};
rpc GetSELinuxInfo(GetSELinuxInfoRequest) returns (GetSELinuxInfoResponse) {};
rpc GetSysConfig(GetSysConfigRequest) returns (GetSysConfigResponse) {};
rpc GetSysErrors(GetSysErrorsRequest) returns (GetSysErrorsResponse) {};
rpc GetMemInfo(GetMemInfoRequest) returns (GetMemInfoResponse) {};
rpc GetMetrics(GetMetricsRequest) returns (GetMetricsResponse) {};
rpc GetProcInfo(GetProcInfoRequest) returns (GetProcInfoResponse) {};
rpc StartProfiling(StartProfilingRequest) returns (StartProfilingResponse) {};
rpc DownloadProfileData(DownloadProfileDataRequest) returns (DownloadProfileDataResponse) {};
rpc GetBucketStats(GetBucketStatsDataRequest) returns (GetBucketStatsDataResponse) {};
rpc GetSRMetrics(GetSRMetricsDataRequest) returns (GetSRMetricsDataResponse) {};
rpc GetAllBucketStats(GetAllBucketStatsRequest) returns (GetAllBucketStatsResponse) {};
rpc LoadBucketMetadata(LoadBucketMetadataRequest) returns (LoadBucketMetadataResponse) {};
rpc DeleteBucketMetadata(DeleteBucketMetadataRequest) returns (DeleteBucketMetadataResponse) {};
rpc DeletePolicy(DeletePolicyRequest) returns (DeletePolicyResponse) {};
rpc LoadPolicy(LoadPolicyRequest) returns (LoadPolicyResponse) {};
rpc LoadPolicyMapping(LoadPolicyMappingRequest) returns (LoadPolicyMappingResponse) {};
rpc DeleteUser(DeleteUserRequest) returns (DeleteUserResponse) {};
rpc DeleteServiceAccount(DeleteServiceAccountRequest) returns (DeleteServiceAccountResponse) {};
rpc LoadUser(LoadUserRequest) returns (LoadUserResponse) {};
rpc LoadServiceAccount(LoadServiceAccountRequest) returns (LoadServiceAccountResponse) {};
rpc LoadGroup(LoadGroupRequest) returns (LoadGroupResponse) {};
rpc ReloadSiteReplicationConfig(ReloadSiteReplicationConfigRequest) returns (ReloadSiteReplicationConfigResponse) {};
// rpc VerifyBinary() returns () {};
// rpc CommitBinary() returns () {};
rpc SignalService(SignalServiceRequest) returns (SignalServiceResponse) {};
rpc BackgroundHealStatus(BackgroundHealStatusRequest) returns (BackgroundHealStatusResponse) {};
rpc GetMetacacheListing(GetMetacacheListingRequest) returns (GetMetacacheListingResponse) {};
rpc UpdateMetacacheListing(UpdateMetacacheListingRequest) returns (UpdateMetacacheListingResponse) {};
rpc ReloadPoolMeta(ReloadPoolMetaRequest) returns (ReloadPoolMetaResponse) {};
rpc StopRebalance(StopRebalanceRequest) returns (StopRebalanceResponse) {};
rpc LoadRebalanceMeta(LoadRebalanceMetaRequest) returns (LoadRebalanceMetaResponse) {};
rpc LoadTransitionTierConfig(LoadTransitionTierConfigRequest) returns (LoadTransitionTierConfigResponse) {};
}
+2
View File
@@ -14,6 +14,8 @@ flatbuffers.workspace = true
lazy_static.workspace = true
lock.workspace = true
protos.workspace = true
rmp-serde.workspace = true
serde.workspace = true
serde_json.workspace = true
tonic = { version = "0.12.3", features = ["gzip"] }
tokio = { workspace = true }
+25 -3
View File
@@ -1,12 +1,16 @@
#![cfg(test)]
use ecstore::disk::VolumeInfo;
use ecstore::{disk::VolumeInfo, store_api::StorageInfo};
use protos::{
models::{PingBody, PingBodyBuilder},
node_service_time_out_client,
proto_gen::node_service::{ListVolumesRequest, MakeVolumeRequest, PingRequest, PingResponse, ReadAllRequest},
proto_gen::node_service::{
ListVolumesRequest, LocalStorageInfoRequest, MakeVolumeRequest, PingRequest, PingResponse, ReadAllRequest,
},
};
use std::error::Error;
use rmp_serde::Deserializer;
use serde::Deserialize;
use std::{error::Error, io::Cursor};
use tonic::Request;
const CLUSTER_ADDR: &str = "http://localhost:9000";
@@ -100,3 +104,21 @@ async fn read_all() -> Result<(), Box<dyn Error>> {
println!("{:?}", volume_infos);
Ok(())
}
#[tokio::test]
async fn storage_info() -> Result<(), Box<dyn Error>> {
let mut client = node_service_time_out_client(&CLUSTER_ADDR.to_string()).await?;
let request = Request::new(LocalStorageInfoRequest { metrics: true });
let response = client.local_storage_info(request).await?.into_inner();
if !response.success {
println!("{:?}", response.error_info);
return Ok(());
}
let info = response.storage_info;
let mut buf = Deserializer::new(Cursor::new(info));
let storage_info: StorageInfo = Deserialize::deserialize(&mut buf).unwrap();
println!("{:?}", storage_info);
Ok(())
}
+3 -2
View File
@@ -17,6 +17,7 @@ common.workspace = true
reader.workspace = true
glob = "0.3.1"
thiserror.workspace = true
flatbuffers.workspace = true
futures.workspace = true
tracing.workspace = true
serde.workspace = true
@@ -38,7 +39,8 @@ netif = "0.1.6"
nix = { version = "0.29.0", features = ["fs"] }
path-absolutize = "3.1.1"
protos.workspace = true
rmp-serde = "1.3.0"
rmp.workspace = true
rmp-serde.workspace = true
tokio-util = { version = "0.7.12", features = ["io", "compat"] }
crc32fast = "1.4.2"
siphasher = "1.0.1"
@@ -51,7 +53,6 @@ tokio = { workspace = true, features = ["io-util", "sync", "signal"] }
tokio-stream = "0.1.15"
tonic.workspace = true
tower.workspace = true
rmp = "0.8.14"
byteorder = "1.5.0"
xxhash-rust = { version = "0.8.12", features = ["xxh64"] }
num = "0.4.3"
+173
View File
@@ -0,0 +1,173 @@
use std::{
collections::{HashMap, HashSet},
time::{SystemTime, UNIX_EPOCH},
};
use common::{
error::{Error, Result},
globals::GLOBAL_Local_Node_Name,
};
use protos::{
models::{PingBody, PingBodyBuilder},
node_service_time_out_client,
proto_gen::node_service::{PingRequest, PingResponse},
};
use serde::{Deserialize, Serialize};
use tonic::Request;
use crate::{
disk::endpoint::Endpoint,
global::GLOBAL_Endpoints,
new_object_layer_fn,
store_api::{StorageAPI, StorageDisk},
};
pub const ITEM_OFFLINE: &str = "offline";
pub const ITEM_INITIALIZING: &str = "initializing";
pub const ITEM_ONLINE: &str = "online";
#[derive(Debug, Default, Serialize, Deserialize)]
pub struct MemStats {
alloc: u64,
total_alloc: u64,
mallocs: u64,
frees: u64,
heap_alloc: u64,
}
#[derive(Debug, Default, Serialize, Deserialize)]
pub struct ServerProperties {
state: String,
endpoint: String,
scheme: String,
uptime: u64,
version: String,
commit_id: String,
network: HashMap<String, String>,
disks: Vec<StorageDisk>,
pool_number: i32,
pool_numbers: Vec<i32>,
mem_stats: MemStats,
max_procs: u64,
num_cpu: u64,
runtime_version: String,
rustfs_env_vars: HashMap<String, String>,
}
async fn is_server_resolvable(endpoint: &Endpoint) -> Result<()> {
let addr = format!(
"{}://{}:{}",
endpoint.url.scheme(),
endpoint.url.host_str().unwrap(),
endpoint.url.port().unwrap()
);
let mut fbb = flatbuffers::FlatBufferBuilder::new();
let payload = fbb.create_vector(b"hello world");
let mut builder = PingBodyBuilder::new(&mut fbb);
builder.add_payload(payload);
let root = builder.finish();
fbb.finish(root, None);
let finished_data = fbb.finished_data();
let decoded_payload = flatbuffers::root::<PingBody>(finished_data);
assert!(decoded_payload.is_ok());
// 创建客户端
let mut client = node_service_time_out_client(&addr)
.await
.map_err(|err| Error::msg(err.to_string()))?;
// 构造 PingRequest
let request = Request::new(PingRequest {
version: 1,
body: finished_data.to_vec(),
});
// 发送请求并获取响应
let response: PingResponse = client.ping(request).await?.into_inner();
// 打印响应
let ping_response_body = flatbuffers::root::<PingBody>(&response.body);
if let Err(e) = ping_response_body {
eprintln!("{}", e);
} else {
println!("ping_resp:body(flatbuffer): {:?}", ping_response_body);
}
Ok(())
}
pub async fn get_local_server_property() -> ServerProperties {
let addr = GLOBAL_Local_Node_Name.read().await.clone();
let mut pool_numbers = HashSet::new();
let mut network = HashMap::new();
let endpoints = match GLOBAL_Endpoints.get() {
Some(eps) => eps,
None => return ServerProperties::default(),
};
for ep in endpoints.as_ref().iter() {
for endpoint in ep.endpoints.as_ref().iter() {
let node_name = match endpoint.url.host_str() {
Some(s) => s.to_string(),
None => addr.clone(),
};
if endpoint.is_local {
pool_numbers.insert(endpoint.pool_idx + 1);
network.insert(node_name, ITEM_ONLINE.to_string());
continue;
}
if !network.contains_key(&node_name) {
if is_server_resolvable(endpoint).await.is_err() {
network.insert(node_name, ITEM_OFFLINE.to_string());
} else {
network.insert(node_name, ITEM_ONLINE.to_string());
}
}
}
}
// todo: mem collect
// let mem_stats =
let mut props = ServerProperties {
endpoint: addr,
uptime: SystemTime::now().duration_since(UNIX_EPOCH).unwrap().as_secs(),
network,
..Default::default()
};
for pool_num in pool_numbers.iter() {
props.pool_numbers.push(*pool_num);
}
props.pool_numbers.sort();
props.pool_number = if props.pool_numbers.len() == 1 {
props.pool_numbers[1]
} else {
i32::MAX
};
// let mut sensitive = HashSet::new();
// sensitive.insert(ENV_ACCESS_KEY.to_string());
// sensitive.insert(ENV_SECRET_KEY.to_string());
// sensitive.insert(ENV_ROOT_USER.to_string());
// sensitive.insert(ENV_ROOT_PASSWORD.to_string());
let layer = new_object_layer_fn();
let lock = layer.read().await;
match lock.as_ref() {
Some(store) => {
let storage_info = store.local_storage_info().await;
props.state = ITEM_ONLINE.to_string();
props.disks = storage_info.disks;
}
None => {
props.state = ITEM_INITIALIZING.to_string();
// todo: get_offline_disks
// props.disks =
}
};
props
}
+5
View File
@@ -19,6 +19,11 @@ lazy_static! {
pub static ref GLOBAL_ConfigSys: ConfigSys = ConfigSys::new();
}
pub const ENV_ACCESS_KEY: &str = "RUSTFS_ACCESS_KEY";
pub const ENV_SECRET_KEY: &str = "RUSTFS_SECRET_KEY";
pub const ENV_ROOT_USER: &str = "RUSTFS_ROOT_USER";
pub const ENV_ROOT_PASSWORD: &str = "RUSTFS_ROOT_PASSWORD";
pub static RUSTFS_CONFIG_PREFIX: &str = "config";
pub struct ConfigSys {}
+1 -7
View File
@@ -30,13 +30,7 @@ pub struct Endpoint {
impl Display for Endpoint {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
if self.url.scheme() == "file" {
write!(
f,
"{}",
fs::canonicalize(self.url.path())
.map_err(|_| std::fmt::Error)?
.to_string_lossy()
)
write!(f, "{}", self.url.path())
} else {
write!(f, "{}", self.url)
}
+1 -1
View File
@@ -65,7 +65,7 @@ async fn init_background_healing() {
.await;
}
async fn get_local_disks_to_heal() -> Vec<Endpoint> {
pub async fn get_local_disks_to_heal() -> Vec<Endpoint> {
let mut disks_to_heal = Vec::new();
for (_, disk) in GLOBAL_LOCAL_DISK_MAP.read().await.iter() {
if let Some(disk) = disk {
+1
View File
@@ -97,6 +97,7 @@ impl Default for HealStartSuccess {
pub type HealStopSuccess = HealStartSuccess;
#[derive(Debug, Default)]
pub struct HealingDisk {
pub id: String,
pub heal_id: String,
+4 -4
View File
@@ -190,19 +190,19 @@ impl HealSequence {
}
impl HealSequence {
fn _get_scanned_items_count(&self) -> usize {
pub fn get_scanned_items_count(&self) -> usize {
self.scanned_items_map.values().sum()
}
fn _get_scanned_items_map(&self) -> ItemsMap {
pub fn _get_scanned_items_map(&self) -> ItemsMap {
self.scanned_items_map.clone()
}
fn _get_healed_items_map(&self) -> ItemsMap {
pub fn _get_healed_items_map(&self) -> ItemsMap {
self.healed_items_map.clone()
}
fn _get_heal_failed_items_map(&self) -> ItemsMap {
pub fn _get_heal_failed_items_map(&self) -> ItemsMap {
self.heal_failed_items_map.clone()
}
+2 -1
View File
@@ -1,3 +1,4 @@
pub mod admin_server_info;
pub mod bitrot;
pub mod cache_value;
mod chunk_stream;
@@ -8,7 +9,7 @@ pub mod endpoints;
pub mod erasure;
pub mod error;
mod file_meta;
mod global;
pub mod global;
pub mod heal;
pub mod peer;
mod quorum;
+104 -6
View File
@@ -1,12 +1,22 @@
use crate::config::common::{read_config, save_config};
use crate::config::error::ConfigError;
use crate::error::{Error, Result};
use crate::new_object_layer_fn;
use crate::store_api::{StorageAPI, StorageDisk, StorageInfo};
use crate::store_err::StorageError;
use crate::{sets::Sets, store::ECStore};
use serde::Serialize;
use byteorder::{ByteOrder, LittleEndian, WriteBytesExt};
use rmp_serde::{Deserializer, Serializer};
use serde::{Deserialize, Serialize};
use std::io::{Cursor, Write};
use std::sync::Arc;
use time::OffsetDateTime;
#[derive(Debug, Clone, Serialize)]
pub const POOL_META_NAME: &str = "pool.bin";
pub const POOL_META_FORMAT: u16 = 1;
pub const POOL_META_VERSION: u16 = 1;
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct PoolStatus {
pub id: usize,
pub cmd_line: String,
@@ -14,8 +24,9 @@ pub struct PoolStatus {
pub decommission: Option<PoolDecommissionInfo>,
}
#[derive(Debug, Clone)]
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct PoolMeta {
pub version: u16,
pub pools: Vec<PoolStatus>,
pub dont_save: bool,
}
@@ -33,6 +44,7 @@ impl PoolMeta {
}
Self {
version: POOL_META_VERSION,
pools: status,
dont_save: false,
}
@@ -46,6 +58,62 @@ impl PoolMeta {
self.pools[idx].decommission.is_some()
}
pub async fn load(&mut self, store: &ECStore) -> Result<()> {
let data = match read_config(store, POOL_META_NAME).await {
Ok(data) => {
if data.is_empty() {
return Ok(());
} else if data.len() <= 4 {
return Err(Error::from_string("poolMeta: no data"));
}
data
}
Err(err) => {
if let Some(ConfigError::NotFound) = err.downcast_ref::<ConfigError>() {
return Ok(());
}
return Err(err);
}
};
let format = LittleEndian::read_u16(&data[0..2]);
if format != POOL_META_FORMAT {
return Err(Error::msg(format!("PoolMeta: unknown format: {}", format)));
}
let version = LittleEndian::read_u16(&data[2..4]);
if version != POOL_META_VERSION {
return Err(Error::msg(format!("PoolMeta: unknown version: {}", version)));
}
let mut buf = Deserializer::new(Cursor::new(&data[4..]));
let meta: PoolMeta = Deserialize::deserialize(&mut buf).unwrap();
*self = meta;
if self.version != POOL_META_VERSION {
return Err(Error::msg(format!("unexpected PoolMeta version: {}", self.version)));
}
Ok(())
}
pub async fn save(&self) -> Result<()> {
if self.dont_save {
return Ok(());
}
let mut data = Vec::new();
data.write_u16::<LittleEndian>(POOL_META_FORMAT).unwrap();
data.write_u16::<LittleEndian>(POOL_META_VERSION).unwrap();
let mut buf = Vec::new();
self.serialize(&mut Serializer::new(&mut buf))?;
data.write_all(&buf)?;
let layer = new_object_layer_fn();
let lock = layer.read().await;
let store = match lock.as_ref() {
Some(s) => s,
None => return Err(Error::from_string("errServerNotInitialized".to_string())),
};
save_config(store, &POOL_META_NAME, &data).await
}
pub fn decommission_cancel(&mut self, idx: usize) -> bool {
if let Some(stats) = self.pools.get_mut(idx) {
if let Some(d) = &stats.decommission {
@@ -72,7 +140,7 @@ impl PoolMeta {
}
}
#[derive(Debug, Clone, Serialize, Default)]
#[derive(Debug, Clone, Serialize, Deserialize, Default)]
pub struct PoolDecommissionInfo {
pub start_time: Option<OffsetDateTime>,
pub start_size: usize,
@@ -103,7 +171,7 @@ pub struct PoolSpaceInfo {
impl ECStore {
pub async fn status(&self, idx: usize) -> Result<PoolStatus> {
let space_info = self.get_decommission_pool_space_info(idx).await?;
let mut pool_info = self.pool_meta.pools[idx].clone();
let mut pool_info = self.pool_meta.read().unwrap().pools[idx].clone();
if let Some(d) = pool_info.decommission.as_mut() {
d.total_size = space_info.total;
d.current_size = space_info.free;
@@ -149,7 +217,7 @@ impl ECStore {
return Err(Error::new(StorageError::DecommissionNotStarted));
}
if self.pool_meta.decommission_cancel(idx) {
if self.pool_meta.write().unwrap().decommission_cancel(idx) {
// FIXME:
}
@@ -182,3 +250,33 @@ fn get_total_usable_capacity_free(disks: &Vec<StorageDisk>, info: &StorageInfo)
}
capacity
}
#[test]
fn test_pool_meta() -> Result<()> {
let meta = PoolMeta::new(vec![]);
let mut data = Vec::new();
data.write_u16::<LittleEndian>(POOL_META_FORMAT).unwrap();
data.write_u16::<LittleEndian>(POOL_META_VERSION).unwrap();
let mut buf = Vec::new();
meta.serialize(&mut Serializer::new(&mut buf))?;
data.write_all(&buf)?;
let format = LittleEndian::read_u16(&data[0..2]);
if format != POOL_META_FORMAT {
return Err(Error::msg(format!("PoolMeta: unknown format: {}", format)));
}
let version = LittleEndian::read_u16(&data[2..4]);
if version != POOL_META_VERSION {
return Err(Error::msg(format!("PoolMeta: unknown version: {}", version)));
}
let mut buf = Deserializer::new(Cursor::new(&data[4..]));
let de_meta: PoolMeta = Deserialize::deserialize(&mut buf).unwrap();
if de_meta.version != POOL_META_VERSION {
return Err(Error::msg(format!("unexpected PoolMeta version: {}", de_meta.version)));
}
println!("meta: {:?}", de_meta);
Ok(())
}
+36 -9
View File
@@ -55,7 +55,7 @@ use std::slice::Iter;
use std::time::SystemTime;
use std::{
collections::{HashMap, HashSet},
sync::Arc,
sync::{Arc, RwLock as std_RwLock},
time::Duration,
};
use time::OffsetDateTime;
@@ -68,7 +68,7 @@ use uuid::Uuid;
const MAX_UPLOADS_LIST: usize = 10000;
#[derive(Debug, Clone)]
#[derive(Debug)]
pub struct ECStore {
pub id: uuid::Uuid,
// pub disks: Vec<DiskStore>,
@@ -76,10 +76,27 @@ pub struct ECStore {
pub pools: Vec<Arc<Sets>>,
pub peer_sys: S3PeerSys,
// pub local_disks: Vec<DiskStore>,
pub pool_meta: PoolMeta,
pub pool_meta: std_RwLock<PoolMeta>,
pub decommission_cancelers: Vec<Option<usize>>,
}
impl Clone for ECStore {
fn clone(&self) -> Self {
let pool_meta = match self.pool_meta.read() {
Ok(pool_meta) => pool_meta.clone(),
Err(_) => PoolMeta::default(),
};
Self {
id: self.id.clone(),
disk_map: self.disk_map.clone(),
pools: self.pools.clone(),
peer_sys: self.peer_sys.clone(),
pool_meta: std_RwLock::new(pool_meta),
decommission_cancelers: self.decommission_cancelers.clone(),
}
}
}
impl ECStore {
#[allow(clippy::new_ret_no_self)]
pub async fn new(_address: String, endpoint_pools: EndpointServerPools) -> Result<Self> {
@@ -189,7 +206,7 @@ impl ECStore {
if !is_dist_erasure().await {
let mut global_local_disk_map = GLOBAL_LOCAL_DISK_MAP.write().await;
for disk in local_disks {
let path = disk.path().to_string_lossy().to_string();
let path = disk.endpoint().to_string();
global_local_disk_map.insert(path, Some(disk.clone()));
}
}
@@ -204,7 +221,7 @@ impl ECStore {
disk_map,
pools,
peer_sys,
pool_meta,
pool_meta: pool_meta.into(),
decommission_cancelers,
};
@@ -480,7 +497,10 @@ impl ECStore {
fn is_suspended(&self, idx: usize) -> bool {
// TODO: LOCK
self.pool_meta.is_suspended(idx)
match self.pool_meta.read() {
Ok(pool_meta) => pool_meta.is_suspended(idx),
Err(_) => false,
}
}
async fn get_pool_idx(&self, bucket: &str, object: &str, size: i64) -> Result<usize> {
@@ -575,7 +595,7 @@ impl ECStore {
let mut has_def_pool = false;
for pinfo in ress.iter() {
if opts.skip_decommissioned && self.pool_meta.is_suspended(pinfo.index) {
if opts.skip_decommissioned && self.pool_meta.read().unwrap().is_suspended(pinfo.index) {
continue;
}
@@ -616,7 +636,7 @@ impl ECStore {
fn pools_with_object(&self, pools: &Vec<PoolObjInfo>, opts: &ObjectOptions) -> Vec<PoolErr> {
let mut errs = Vec::new();
for pool in pools.iter() {
if opts.skip_decommissioned && self.pool_meta.is_suspended(pool.index) {
if opts.skip_decommissioned && self.pool_meta.read().unwrap().is_suspended(pool.index) {
continue;
}
// TODO:SkipRebalancing
@@ -869,6 +889,13 @@ impl ECStore {
Ok(objs[0].as_ref().unwrap().clone())
}
pub async fn reload_pool_meta(&self) -> Result<()> {
let mut meta = PoolMeta::default();
meta.load(self).await?;
*self.pool_meta.write().unwrap() = meta;
Ok(())
}
}
async fn update_scan(
@@ -965,7 +992,7 @@ pub async fn init_local_disks(endpoint_pools: EndpointServerPools) -> Result<()>
let disk = new_disk(ep, opt).await?;
let path = disk.path().to_string_lossy().to_string();
let path = disk.endpoint().to_string();
global_local_disk_map.insert(path, Some(disk.clone()));
+4 -4
View File
@@ -795,7 +795,7 @@ pub struct DeletedObject {
// pub replication_state: ReplicationState,
}
#[derive(Debug, Default, Serialize)]
#[derive(Debug, Default, Serialize, Deserialize)]
pub enum BackendByte {
#[default]
Unknown,
@@ -833,13 +833,13 @@ pub struct StorageDisk {
pub disk_index: i32,
}
#[derive(Debug, Default)]
#[derive(Debug, Default, Serialize, Deserialize)]
pub struct StorageInfo {
pub disks: Vec<StorageDisk>,
pub backend: BackendInfo,
}
#[derive(Debug, Default, Serialize)]
#[derive(Debug, Default, Serialize, Deserialize)]
pub struct BackendDisks(HashMap<String, usize>);
impl BackendDisks {
@@ -851,7 +851,7 @@ impl BackendDisks {
}
}
#[derive(Debug, Default, Serialize)]
#[derive(Debug, Default, Serialize, Deserialize)]
#[serde(rename_all = "PascalCase", default)]
pub struct BackendInfo {
pub backend_type: BackendByte,
+12
View File
@@ -0,0 +1,12 @@
[package]
name = "madmin"
edition.workspace = true
license.workspace = true
repository.workspace = true
rust-version.workspace = true
version.workspace = true
[dependencies]
ecstore.workspace = true
psutil = "3.3.0"
serde.workspace = true
+115
View File
@@ -0,0 +1,115 @@
use std::collections::{HashMap, HashSet};
use ecstore::{
config::storageclass::{RRS, STANDARD},
global::GLOBAL_BackgroundHealState,
heal::{background_heal_ops::get_local_disks_to_heal, heal_ops::BG_HEALING_UUID},
new_object_layer_fn,
store_api::{StorageAPI, StorageDisk},
};
use serde::{Deserialize, Serialize};
#[derive(Debug, Default, Serialize, Deserialize)]
pub struct MRFStatus {
bytes_healed: u64,
items_healed: u64,
}
#[derive(Debug, Default, Serialize, Deserialize)]
pub struct SetStatus {
pub id: String,
pub pool_index: i32,
pub set_index: i32,
pub heal_status: String,
pub heal_priority: String,
pub total_objects: usize,
pub disks: Vec<StorageDisk>,
}
#[derive(Debug, Default, Serialize, Deserialize)]
pub struct BgHealState {
offline_endpoints: Vec<String>,
scanned_items_count: u64,
heal_disks: Vec<String>,
sets: Vec<SetStatus>,
mrf: HashMap<String, MRFStatus>,
scparity: HashMap<String, usize>,
}
pub async fn get_local_background_heal_status() -> (BgHealState, bool) {
let (bg_seq, ok) = GLOBAL_BackgroundHealState
.read()
.await
.get_heal_sequence_by_token(BG_HEALING_UUID)
.await;
if !ok {
return (BgHealState::default(), false);
}
let bg_seq = bg_seq.unwrap();
let mut status = BgHealState {
scanned_items_count: bg_seq.read().await.get_scanned_items_count() as u64,
..Default::default()
};
let mut heal_disks_map = HashSet::new();
for ep in get_local_disks_to_heal().await.iter() {
heal_disks_map.insert(ep.to_string());
}
let layer = new_object_layer_fn();
let lock = layer.read().await;
let store = match lock.as_ref() {
Some(s) => s,
None => {
let healing = GLOBAL_BackgroundHealState.read().await.get_local_healing_disks().await;
for disk in healing.values() {
status.heal_disks.push(disk.endpoint.clone());
}
return (status, true);
}
};
let si = store.local_storage_info().await;
let mut indexed = HashMap::new();
for disk in si.disks.iter() {
let set_idx = format!("{}-{}", disk.pool_index, disk.set_index);
// indexed.insert(set_idx, disk);
indexed.entry(set_idx).or_insert(Vec::new()).push(disk);
}
for (id, disks) in indexed {
let mut ss = SetStatus {
id,
set_index: disks[0].set_index,
pool_index: disks[0].pool_index,
..Default::default()
};
for disk in disks {
ss.disks.push(disk.clone());
if disk.healing {
ss.heal_status = "healing".to_string();
ss.heal_priority = "high".to_string();
status.heal_disks.push(disk.endpoint.clone());
}
}
ss.disks.sort_by(|a, b| {
if a.pool_index != b.pool_index {
return a.pool_index.cmp(&b.pool_index);
}
if a.set_index != b.set_index {
return a.set_index.cmp(&b.set_index);
}
a.disk_index.cmp(&b.disk_index)
});
status.sets.push(ss);
}
status.sets.sort_by(|a, b| a.id.cmp(&b.id));
let backend_info = store.backend_info().await;
status
.scparity
.insert(STANDARD.to_string(), backend_info.standard_sc_parity.unwrap_or_default());
status
.scparity
.insert(RRS.to_string(), backend_info.rr_sc_parity.unwrap_or_default());
(status, true)
}
+177
View File
@@ -0,0 +1,177 @@
use std::collections::HashMap;
use serde::{Deserialize, Serialize};
#[derive(Debug, Default, Serialize, Deserialize)]
pub struct NodeCommon {
pub addr: String,
pub error: String,
}
#[derive(Debug, Default, Serialize, Deserialize)]
pub struct Cpu {
vendor_id: String,
family: String,
model: String,
stepping: i32,
physical_id: String,
model_name: String,
mhz: f64,
cache_size: i32,
flags: Vec<String>,
microcode: String,
cores: u64,
}
#[derive(Debug, Default, Serialize, Deserialize)]
pub struct CpuFreqStats {
name: String,
cpuinfo_current_frequency: Option<u64>,
cpuinfo_minimum_frequency: Option<u64>,
cpuinfo_maximum_frequency: Option<u64>,
cpuinfo_transition_latency: Option<u64>,
scaling_current_frequency: Option<u64>,
scaling_minimum_frequency: Option<u64>,
scaling_maximum_frequency: Option<u64>,
available_governors: String,
driver: String,
governor: String,
related_cpus: String,
set_speed: String,
}
#[derive(Debug, Default, Serialize, Deserialize)]
pub struct Cpus {
node_common: NodeCommon,
cpus: Vec<Cpu>,
cpu_freq_stats: Vec<CpuFreqStats>,
}
pub fn get_cpus() -> Cpus {
// todo
Cpus::default()
}
#[derive(Debug, Default, Serialize, Deserialize)]
pub struct Partition {
pub error: String,
device: String,
model: String,
revision: String,
mountpoint: String,
fs_type: String,
mount_options: String,
space_total: u64,
space_free: u64,
inode_total: u64,
inode_free: u64,
}
#[derive(Debug, Default, Serialize, Deserialize)]
pub struct Partitions {
node_common: NodeCommon,
partitions: Vec<Partition>,
}
pub fn get_partitions() -> Partitions {
Partitions::default()
}
#[derive(Debug, Default, Serialize, Deserialize)]
pub struct OsInfo {
node_common: NodeCommon,
}
pub fn get_os_info() -> OsInfo {
OsInfo::default()
}
#[derive(Debug, Default, Serialize, Deserialize)]
pub struct ProcInfo {
node_common: NodeCommon,
pid: i32,
is_background: bool,
cpu_percent: f64,
children_pids: Vec<i32>,
cmd_line: String,
num_connections: usize,
create_time: u64,
cwd: String,
exec_path: String,
gids: Vec<i32>,
// io_counters:
is_running: bool,
// mem_info:
// mem_maps:
mem_percent: f32,
name: String,
nice: i32,
//num_ctx_switches:
num_fds: i32,
num_threads: i32,
// page_faults:
ppid: i32,
status: String,
tgid: i32,
uids: Vec<i32>,
username: String,
}
pub fn get_proc_info(_addr: &str) -> ProcInfo {
ProcInfo::default()
}
#[derive(Debug, Default, Serialize, Deserialize)]
pub struct SysService {
name: String,
status: String,
}
#[derive(Debug, Default, Serialize, Deserialize)]
pub struct SysServices {
node_common: NodeCommon,
services: Vec<SysService>,
}
pub fn get_sys_services(_add: &str) -> SysServices {
SysServices::default()
}
#[derive(Debug, Default, Serialize, Deserialize)]
pub struct SysConfig {
node_common: NodeCommon,
config: HashMap<String, String>,
}
pub fn get_sys_config(_addr: &str) -> SysConfig {
SysConfig::default()
}
#[derive(Debug, Default, Serialize, Deserialize)]
pub struct SysErrors {
node_common: NodeCommon,
errors: Vec<String>,
}
pub fn get_sys_errors(_add: &str) -> SysErrors {
SysErrors::default()
}
#[derive(Debug, Default, Serialize, Deserialize)]
pub struct MemInfo {
node_common: NodeCommon,
total: u64,
used: u64,
free: u64,
available: u64,
shared: u64,
cache: u64,
buffers: u64,
swap_space_total: u64,
swap_space_free: u64,
limit: u64,
}
pub fn get_mem_info(_addr: &str) -> MemInfo {
MemInfo::default()
}
+4
View File
@@ -0,0 +1,4 @@
pub mod heal_command;
pub mod health;
pub mod metrics;
pub mod net;
+69
View File
@@ -0,0 +1,69 @@
use std::collections::{HashMap, HashSet};
use serde::{Deserialize, Serialize};
#[derive(Debug, Default, Serialize, Deserialize)]
pub struct TimedAction {
count: u64,
acc_time: u64,
bytes: u64,
}
#[derive(Debug, Default, Serialize, Deserialize)]
pub struct DiskIOStats {
read_ios: u64,
read_merges: u64,
read_sectors: u64,
read_ticks: u64,
write_ios: u64,
write_merges: u64,
write_sectors: u64,
write_ticks: u64,
current_ios: u64,
total_ticks: u64,
req_ticks: u64,
discard_ios: u64,
discard_merges: u64,
discard_sectors: u64,
discard_ticks: u64,
flush_ios: u64,
flush_ticks: u64,
}
#[derive(Debug, Default, Serialize, Deserialize)]
pub struct DiskMetric {
collected_at: u64,
n_disks: usize,
offline: usize,
healing: usize,
life_time_ops: HashMap<String, u64>,
last_minute: HashMap<String, TimedAction>,
io_stats: DiskIOStats,
}
#[derive(Debug, Default, Serialize, Deserialize)]
pub struct Metrics {}
#[derive(Debug, Default, Serialize, Deserialize)]
pub struct RealtimeMetrics {
errors: Vec<String>,
hosts: Vec<String>,
aggregated: Metrics,
by_host: HashMap<String, Metrics>,
by_disk: HashMap<String, DiskMetric>,
finally: bool,
}
#[derive(Debug, Default, Serialize, Deserialize)]
pub struct CollectMetricsOpts {
hosts: HashSet<String>,
disks: HashSet<String>,
job_id: String,
dep_id: String,
}
pub type MetricType = u64;
pub fn collect_local_metrics(_types: MetricType, _opts: &CollectMetricsOpts) -> RealtimeMetrics {
RealtimeMetrics::default()
}
+14
View File
@@ -0,0 +1,14 @@
use serde::{Deserialize, Serialize};
use crate::health::NodeCommon;
#[cfg(target_os = "linux")]
pub mod net_linux;
#[derive(Debug, Default, Serialize, Deserialize)]
pub struct NetInfo {
node_common: NodeCommon,
interface: String,
driver: String,
firmware_version: String,
}
+9
View File
@@ -0,0 +1,9 @@
use super::NetInfo;
pub fn get_net_info(addr: &str, iface: &str) -> NetInfo {
let mut ni = NetInfo::default();
ni.node_common.addr = addr.to_string();
ni.interface = iface.to_string();
ni
}
+4 -1
View File
@@ -9,6 +9,7 @@ rust-version.workspace = true
# See more keys and their definitions at https://doc.rust-lang.org/cargo/reference/manifest.html
[dependencies]
madmin.workspace = true
log.workspace = true
async-trait.workspace = true
bytes.workspace = true
@@ -31,7 +32,9 @@ prost.workspace = true
prost-types.workspace = true
protos.workspace = true
protobuf.workspace = true
rmp-serde.workspace = true
s3s.workspace = true
serde.workspace = true
serde_json.workspace = true
tracing.workspace = true
time = { workspace = true, features = ["parsing", "formatting"] }
@@ -51,13 +54,13 @@ tracing-error.workspace = true
tracing-subscriber.workspace = true
transform-stream.workspace = true
uuid = "1.11.0"
url.workspace = true
admin = { path = "../api/admin" }
axum.workspace = true
matchit = "0.8.4"
shadow-rs = "0.35.2"
const-str = { version = "0.5.7", features = ["std", "proc"] }
atoi = "2.0.0"
serde.workspace = true
serde_urlencoded = "0.7.1"
[build-dependencies]
+676 -16
View File
@@ -1,36 +1,43 @@
use std::{error::Error, io::ErrorKind, pin::Pin};
use std::{
collections::HashMap,
error::Error,
io::{Cursor, ErrorKind},
pin::Pin,
sync::Arc,
};
use ecstore::{
admin_server_info::get_local_server_property,
bucket::{metadata::load_bucket_metadata, metadata_sys::GLOBAL_BucketMetadataSys},
disk::{
DeleteOptions, DiskAPI, DiskInfoOptions, DiskStore, FileInfoVersions, ReadMultipleReq, ReadOptions, Reader,
UpdateMetadataOpts, WalkDirOptions,
},
erasure::Writer,
heal::{data_usage_cache::DataUsageCache, heal_commands::HealOpts},
new_object_layer_fn,
peer::{LocalPeerS3Client, PeerS3Client},
store::{all_local_disk_path, find_local_disk},
store_api::{BucketOptions, DeleteBucketOptions, FileInfo, MakeBucketOptions},
store_api::{BucketOptions, DeleteBucketOptions, FileInfo, MakeBucketOptions, StorageAPI},
};
use futures::{Stream, StreamExt};
use lock::{lock_args::LockArgs, Locker, GLOBAL_LOCAL_SERVER};
use common::globals::GLOBAL_Local_Node_Name;
use madmin::net::net_linux::get_net_info;
use madmin::{
heal_command::get_local_background_heal_status,
health::{
get_cpus, get_mem_info, get_os_info, get_partitions, get_proc_info, get_sys_config, get_sys_errors, get_sys_services,
},
metrics::{collect_local_metrics, CollectMetricsOpts},
};
use protos::{
models::{PingBody, PingBodyBuilder},
proto_gen::node_service::{
node_service_server::NodeService as Node, CheckPartsRequest, CheckPartsResponse, DeleteBucketRequest,
DeleteBucketResponse, DeletePathsRequest, DeletePathsResponse, DeleteRequest, DeleteResponse, DeleteVersionRequest,
DeleteVersionResponse, DeleteVersionsRequest, DeleteVersionsResponse, DeleteVolumeRequest, DeleteVolumeResponse,
DiskInfoRequest, DiskInfoResponse, GenerallyLockRequest, GenerallyLockResponse, GetBucketInfoRequest,
GetBucketInfoResponse, HealBucketRequest, HealBucketResponse, ListBucketRequest, ListBucketResponse, ListDirRequest,
ListDirResponse, ListVolumesRequest, ListVolumesResponse, MakeBucketRequest, MakeBucketResponse, MakeVolumeRequest,
MakeVolumeResponse, MakeVolumesRequest, MakeVolumesResponse, NsScannerRequest, NsScannerResponse, PingRequest,
PingResponse, ReadAllRequest, ReadAllResponse, ReadAtRequest, ReadAtResponse, ReadMultipleRequest, ReadMultipleResponse,
ReadVersionRequest, ReadVersionResponse, ReadXlRequest, ReadXlResponse, RenameDataRequest, RenameDataResponse,
RenameFileRequst, RenameFileResponse, RenamePartRequst, RenamePartResponse, StatVolumeRequest, StatVolumeResponse,
UpdateMetadataRequest, UpdateMetadataResponse, VerifyFileRequest, VerifyFileResponse, WalkDirRequest, WalkDirResponse,
WriteAllRequest, WriteAllResponse, WriteMetadataRequest, WriteMetadataResponse, WriteRequest, WriteResponse,
},
proto_gen::node_service::{node_service_server::NodeService as Node, *},
};
use rmp_serde::{Deserializer, Serializer};
use serde::{Deserialize, Serialize};
use tokio::sync::mpsc;
use tokio_stream::wrappers::ReceiverStream;
use tonic::{Request, Response, Status, Streaming};
@@ -1506,4 +1513,657 @@ impl Node for NodeService {
})),
}
}
async fn local_storage_info(
&self,
_request: Request<LocalStorageInfoRequest>,
) -> Result<Response<LocalStorageInfoResponse>, Status> {
// let request = request.into_inner();
let layer = new_object_layer_fn();
let lock = layer.read().await;
let store = match lock.as_ref() {
Some(s) => s,
None => {
return Ok(tonic::Response::new(LocalStorageInfoResponse {
success: false,
storage_info: vec![],
error_info: Some("errServerNotInitialized".to_string()),
}))
}
};
let info = store.local_storage_info().await;
let mut buf = Vec::new();
if let Err(err) = info.serialize(&mut Serializer::new(&mut buf)) {
return Ok(tonic::Response::new(LocalStorageInfoResponse {
success: false,
storage_info: vec![],
error_info: Some(err.to_string()),
}));
}
Ok(tonic::Response::new(LocalStorageInfoResponse {
success: true,
storage_info: buf,
error_info: None,
}))
}
async fn server_info(&self, _request: Request<ServerInfoRequest>) -> Result<Response<ServerInfoResponse>, Status> {
let info = get_local_server_property().await;
let mut buf = Vec::new();
if let Err(err) = info.serialize(&mut Serializer::new(&mut buf)) {
return Ok(tonic::Response::new(ServerInfoResponse {
success: false,
server_properties: vec![],
error_info: Some(err.to_string()),
}));
}
Ok(tonic::Response::new(ServerInfoResponse {
success: true,
server_properties: buf,
error_info: None,
}))
}
async fn get_cpus(&self, _request: Request<GetCpusRequest>) -> Result<Response<GetCpusResponse>, Status> {
let info = get_cpus();
let mut buf = Vec::new();
if let Err(err) = info.serialize(&mut Serializer::new(&mut buf)) {
return Ok(tonic::Response::new(GetCpusResponse {
success: false,
cpus: vec![],
error_info: Some(err.to_string()),
}));
}
Ok(tonic::Response::new(GetCpusResponse {
success: true,
cpus: buf,
error_info: None,
}))
}
async fn get_net_info(&self, _request: Request<GetNetInfoRequest>) -> Result<Response<GetNetInfoResponse>, Status> {
let addr = GLOBAL_Local_Node_Name.read().await.clone();
let info = get_net_info(&addr, "");
let mut buf = Vec::new();
if let Err(err) = info.serialize(&mut Serializer::new(&mut buf)) {
return Ok(tonic::Response::new(GetNetInfoResponse {
success: false,
net_info: vec![],
error_info: Some(err.to_string()),
}));
}
Ok(tonic::Response::new(GetNetInfoResponse {
success: true,
net_info: buf,
error_info: None,
}))
}
async fn get_partitions(&self, _request: Request<GetPartitionsRequest>) -> Result<Response<GetPartitionsResponse>, Status> {
let partitions = get_partitions();
let mut buf = Vec::new();
if let Err(err) = partitions.serialize(&mut Serializer::new(&mut buf)) {
return Ok(tonic::Response::new(GetPartitionsResponse {
success: false,
partitions: vec![],
error_info: Some(err.to_string()),
}));
}
Ok(tonic::Response::new(GetPartitionsResponse {
success: true,
partitions: buf,
error_info: None,
}))
}
async fn get_os_info(&self, _request: Request<GetOsInfoRequest>) -> Result<Response<GetOsInfoResponse>, Status> {
let os_info = get_os_info();
let mut buf = Vec::new();
if let Err(err) = os_info.serialize(&mut Serializer::new(&mut buf)) {
return Ok(tonic::Response::new(GetOsInfoResponse {
success: false,
os_info: vec![],
error_info: Some(err.to_string()),
}));
}
Ok(tonic::Response::new(GetOsInfoResponse {
success: true,
os_info: buf,
error_info: None,
}))
}
async fn get_se_linux_info(
&self,
_request: Request<GetSeLinuxInfoRequest>,
) -> Result<Response<GetSeLinuxInfoResponse>, Status> {
let addr = GLOBAL_Local_Node_Name.read().await.clone();
let info = get_sys_services(&addr);
let mut buf = Vec::new();
if let Err(err) = info.serialize(&mut Serializer::new(&mut buf)) {
return Ok(tonic::Response::new(GetSeLinuxInfoResponse {
success: false,
sys_services: vec![],
error_info: Some(err.to_string()),
}));
}
Ok(tonic::Response::new(GetSeLinuxInfoResponse {
success: true,
sys_services: buf,
error_info: None,
}))
}
async fn get_sys_config(&self, _request: Request<GetSysConfigRequest>) -> Result<Response<GetSysConfigResponse>, Status> {
let addr = GLOBAL_Local_Node_Name.read().await.clone();
let info = get_sys_config(&addr);
let mut buf = Vec::new();
if let Err(err) = info.serialize(&mut Serializer::new(&mut buf)) {
return Ok(tonic::Response::new(GetSysConfigResponse {
success: false,
sys_config: vec![],
error_info: Some(err.to_string()),
}));
}
Ok(tonic::Response::new(GetSysConfigResponse {
success: true,
sys_config: buf,
error_info: None,
}))
}
async fn get_sys_errors(&self, _request: Request<GetSysErrorsRequest>) -> Result<Response<GetSysErrorsResponse>, Status> {
let addr = GLOBAL_Local_Node_Name.read().await.clone();
let info = get_sys_errors(&addr);
let mut buf = Vec::new();
if let Err(err) = info.serialize(&mut Serializer::new(&mut buf)) {
return Ok(tonic::Response::new(GetSysErrorsResponse {
success: false,
sys_errors: vec![],
error_info: Some(err.to_string()),
}));
}
Ok(tonic::Response::new(GetSysErrorsResponse {
success: true,
sys_errors: buf,
error_info: None,
}))
}
async fn get_mem_info(&self, _request: Request<GetMemInfoRequest>) -> Result<Response<GetMemInfoResponse>, Status> {
let addr = GLOBAL_Local_Node_Name.read().await.clone();
let info = get_mem_info(&addr);
let mut buf = Vec::new();
if let Err(err) = info.serialize(&mut Serializer::new(&mut buf)) {
return Ok(tonic::Response::new(GetMemInfoResponse {
success: false,
mem_info: vec![],
error_info: Some(err.to_string()),
}));
}
Ok(tonic::Response::new(GetMemInfoResponse {
success: true,
mem_info: buf,
error_info: None,
}))
}
async fn get_metrics(&self, request: Request<GetMetricsRequest>) -> Result<Response<GetMetricsResponse>, Status> {
let request = request.into_inner();
let mut buf = Deserializer::new(Cursor::new(request.opts));
let opts: CollectMetricsOpts = Deserialize::deserialize(&mut buf).unwrap();
let info = collect_local_metrics(request.metric_type, &opts);
let mut buf = Vec::new();
if let Err(err) = info.serialize(&mut Serializer::new(&mut buf)) {
return Ok(tonic::Response::new(GetMetricsResponse {
success: false,
realtime_metrics: vec![],
error_info: Some(err.to_string()),
}));
}
Ok(tonic::Response::new(GetMetricsResponse {
success: true,
realtime_metrics: buf,
error_info: None,
}))
}
async fn get_proc_info(&self, _request: Request<GetProcInfoRequest>) -> Result<Response<GetProcInfoResponse>, Status> {
let addr = GLOBAL_Local_Node_Name.read().await.clone();
let info = get_proc_info(&addr);
let mut buf = Vec::new();
if let Err(err) = info.serialize(&mut Serializer::new(&mut buf)) {
return Ok(tonic::Response::new(GetProcInfoResponse {
success: false,
proc_info: vec![],
error_info: Some(err.to_string()),
}));
}
Ok(tonic::Response::new(GetProcInfoResponse {
success: true,
proc_info: buf,
error_info: None,
}))
}
async fn start_profiling(
&self,
_request: Request<StartProfilingRequest>,
) -> Result<Response<StartProfilingResponse>, Status> {
todo!()
}
async fn download_profile_data(
&self,
_request: Request<DownloadProfileDataRequest>,
) -> Result<Response<DownloadProfileDataResponse>, Status> {
todo!()
}
async fn get_bucket_stats(
&self,
_request: Request<GetBucketStatsDataRequest>,
) -> Result<Response<GetBucketStatsDataResponse>, Status> {
todo!()
}
async fn get_sr_metrics(
&self,
_request: Request<GetSrMetricsDataRequest>,
) -> Result<Response<GetSrMetricsDataResponse>, Status> {
todo!()
}
async fn get_all_bucket_stats(
&self,
_request: Request<GetAllBucketStatsRequest>,
) -> Result<Response<GetAllBucketStatsResponse>, Status> {
todo!()
}
async fn load_bucket_metadata(
&self,
request: Request<LoadBucketMetadataRequest>,
) -> Result<Response<LoadBucketMetadataResponse>, Status> {
let request = request.into_inner();
let bucket = request.bucket;
if bucket.is_empty() {
return Ok(tonic::Response::new(LoadBucketMetadataResponse {
success: false,
error_info: Some("bucket name is missing".to_string()),
}));
}
let layer = new_object_layer_fn();
let lock = layer.read().await;
let store = match lock.as_ref() {
Some(s) => s,
None => {
return Ok(tonic::Response::new(LoadBucketMetadataResponse {
success: false,
error_info: Some("errServerNotInitialized".to_string()),
}))
}
};
match load_bucket_metadata(store, &bucket).await {
Ok(meta) => {
GLOBAL_BucketMetadataSys.write().await.set(bucket, Arc::new(meta)).await;
Ok(tonic::Response::new(LoadBucketMetadataResponse {
success: true,
error_info: None,
}))
}
Err(err) => Ok(tonic::Response::new(LoadBucketMetadataResponse {
success: false,
error_info: Some(err.to_string()),
})),
}
}
async fn delete_bucket_metadata(
&self,
request: Request<DeleteBucketMetadataRequest>,
) -> Result<Response<DeleteBucketMetadataResponse>, Status> {
let request = request.into_inner();
let _bucket = request.bucket;
//todo
Ok(tonic::Response::new(DeleteBucketMetadataResponse {
success: true,
error_info: None,
}))
}
async fn delete_policy(&self, request: Request<DeletePolicyRequest>) -> Result<Response<DeletePolicyResponse>, Status> {
let request = request.into_inner();
let policy = request.policy_name;
if policy.is_empty() {
return Ok(tonic::Response::new(DeletePolicyResponse {
success: false,
error_info: Some("policy name is missing".to_string()),
}));
}
let layer = new_object_layer_fn();
let lock = layer.read().await;
let _store = match lock.as_ref() {
Some(s) => s,
None => {
return Ok(tonic::Response::new(DeletePolicyResponse {
success: false,
error_info: Some("errServerNotInitialized".to_string()),
}))
}
};
todo!()
}
async fn load_policy(&self, request: Request<LoadPolicyRequest>) -> Result<Response<LoadPolicyResponse>, Status> {
let request = request.into_inner();
let policy = request.policy_name;
if policy.is_empty() {
return Ok(tonic::Response::new(LoadPolicyResponse {
success: false,
error_info: Some("policy name is missing".to_string()),
}));
}
let layer = new_object_layer_fn();
let lock = layer.read().await;
let _store = match lock.as_ref() {
Some(s) => s,
None => {
return Ok(tonic::Response::new(LoadPolicyResponse {
success: false,
error_info: Some("errServerNotInitialized".to_string()),
}))
}
};
todo!()
}
async fn load_policy_mapping(
&self,
request: Request<LoadPolicyMappingRequest>,
) -> Result<Response<LoadPolicyMappingResponse>, Status> {
let request = request.into_inner();
let user_or_group = request.user_or_group;
if user_or_group.is_empty() {
return Ok(tonic::Response::new(LoadPolicyMappingResponse {
success: false,
error_info: Some("user_or_group name is missing".to_string()),
}));
}
let _user_type = request.user_type;
let _is_group = request.is_group;
let layer = new_object_layer_fn();
let lock = layer.read().await;
let _store = match lock.as_ref() {
Some(s) => s,
None => {
return Ok(tonic::Response::new(LoadPolicyMappingResponse {
success: false,
error_info: Some("errServerNotInitialized".to_string()),
}))
}
};
todo!()
}
async fn delete_user(&self, request: Request<DeleteUserRequest>) -> Result<Response<DeleteUserResponse>, Status> {
let request = request.into_inner();
let access_key = request.access_key;
if access_key.is_empty() {
return Ok(tonic::Response::new(DeleteUserResponse {
success: false,
error_info: Some("access_key name is missing".to_string()),
}));
}
let layer = new_object_layer_fn();
let lock = layer.read().await;
let _store = match lock.as_ref() {
Some(s) => s,
None => {
return Ok(tonic::Response::new(DeleteUserResponse {
success: false,
error_info: Some("errServerNotInitialized".to_string()),
}))
}
};
todo!()
}
async fn delete_service_account(
&self,
request: Request<DeleteServiceAccountRequest>,
) -> Result<Response<DeleteServiceAccountResponse>, Status> {
let request = request.into_inner();
let access_key = request.access_key;
if access_key.is_empty() {
return Ok(tonic::Response::new(DeleteServiceAccountResponse {
success: false,
error_info: Some("access_key name is missing".to_string()),
}));
}
let layer = new_object_layer_fn();
let lock = layer.read().await;
let _store = match lock.as_ref() {
Some(s) => s,
None => {
return Ok(tonic::Response::new(DeleteServiceAccountResponse {
success: false,
error_info: Some("errServerNotInitialized".to_string()),
}))
}
};
todo!()
}
async fn load_user(&self, request: Request<LoadUserRequest>) -> Result<Response<LoadUserResponse>, Status> {
let request = request.into_inner();
let access_key = request.access_key;
let _temp = request.temp;
if access_key.is_empty() {
return Ok(tonic::Response::new(LoadUserResponse {
success: false,
error_info: Some("access_key name is missing".to_string()),
}));
}
let layer = new_object_layer_fn();
let lock = layer.read().await;
let _store = match lock.as_ref() {
Some(s) => s,
None => {
return Ok(tonic::Response::new(LoadUserResponse {
success: false,
error_info: Some("errServerNotInitialized".to_string()),
}))
}
};
todo!()
}
async fn load_service_account(
&self,
request: Request<LoadServiceAccountRequest>,
) -> Result<Response<LoadServiceAccountResponse>, Status> {
let request = request.into_inner();
let access_key = request.access_key;
if access_key.is_empty() {
return Ok(tonic::Response::new(LoadServiceAccountResponse {
success: false,
error_info: Some("access_key name is missing".to_string()),
}));
}
let layer = new_object_layer_fn();
let lock = layer.read().await;
let _store = match lock.as_ref() {
Some(s) => s,
None => {
return Ok(tonic::Response::new(LoadServiceAccountResponse {
success: false,
error_info: Some("errServerNotInitialized".to_string()),
}))
}
};
todo!()
}
async fn load_group(&self, request: Request<LoadGroupRequest>) -> Result<Response<LoadGroupResponse>, Status> {
let request = request.into_inner();
let group = request.group;
if group.is_empty() {
return Ok(tonic::Response::new(LoadGroupResponse {
success: false,
error_info: Some("group name is missing".to_string()),
}));
}
let layer = new_object_layer_fn();
let lock = layer.read().await;
let _store = match lock.as_ref() {
Some(s) => s,
None => {
return Ok(tonic::Response::new(LoadGroupResponse {
success: false,
error_info: Some("errServerNotInitialized".to_string()),
}))
}
};
todo!()
}
async fn reload_site_replication_config(
&self,
_request: Request<ReloadSiteReplicationConfigRequest>,
) -> Result<Response<ReloadSiteReplicationConfigResponse>, Status> {
let layer = new_object_layer_fn();
let lock = layer.read().await;
let _store = match lock.as_ref() {
Some(s) => s,
None => {
return Ok(tonic::Response::new(ReloadSiteReplicationConfigResponse {
success: false,
error_info: Some("errServerNotInitialized".to_string()),
}))
}
};
todo!()
}
async fn signal_service(&self, request: Request<SignalServiceRequest>) -> Result<Response<SignalServiceResponse>, Status> {
let request = request.into_inner();
let _vars = match request.vars {
Some(vars) => vars.value,
None => HashMap::new(),
};
todo!()
}
async fn background_heal_status(
&self,
_request: Request<BackgroundHealStatusRequest>,
) -> Result<Response<BackgroundHealStatusResponse>, Status> {
let (state, ok) = get_local_background_heal_status().await;
if !ok {
return Ok(tonic::Response::new(BackgroundHealStatusResponse {
success: false,
bg_heal_state: vec![],
error_info: Some("errServerNotInitialized".to_string()),
}));
}
let mut buf = Vec::new();
if let Err(err) = state.serialize(&mut Serializer::new(&mut buf)) {
return Ok(tonic::Response::new(BackgroundHealStatusResponse {
success: false,
bg_heal_state: vec![],
error_info: Some(err.to_string()),
}));
}
Ok(tonic::Response::new(BackgroundHealStatusResponse {
success: true,
bg_heal_state: buf,
error_info: None,
}))
}
async fn get_metacache_listing(
&self,
_request: Request<GetMetacacheListingRequest>,
) -> Result<Response<GetMetacacheListingResponse>, Status> {
todo!()
}
async fn update_metacache_listing(
&self,
_request: Request<UpdateMetacacheListingRequest>,
) -> Result<Response<UpdateMetacacheListingResponse>, Status> {
todo!()
}
async fn reload_pool_meta(
&self,
_request: Request<ReloadPoolMetaRequest>,
) -> Result<Response<ReloadPoolMetaResponse>, Status> {
let layer = new_object_layer_fn();
let lock = layer.read().await;
let store = match lock.as_ref() {
Some(s) => s,
None => {
return Ok(tonic::Response::new(ReloadPoolMetaResponse {
success: false,
error_info: Some("errServerNotInitialized".to_string()),
}))
}
};
match store.reload_pool_meta().await {
Ok(_) => Ok(tonic::Response::new(ReloadPoolMetaResponse {
success: true,
error_info: None,
})),
Err(err) => Ok(tonic::Response::new(ReloadPoolMetaResponse {
success: false,
error_info: Some(err.to_string()),
})),
}
}
async fn stop_rebalance(&self, _request: Request<StopRebalanceRequest>) -> Result<Response<StopRebalanceResponse>, Status> {
let layer = new_object_layer_fn();
let lock = layer.read().await;
let _store = match lock.as_ref() {
Some(s) => s,
None => {
return Ok(tonic::Response::new(StopRebalanceResponse {
success: false,
error_info: Some("errServerNotInitialized".to_string()),
}))
}
};
// todo
// store.stop_rebalance().await;
todo!()
}
async fn load_rebalance_meta(
&self,
_request: Request<LoadRebalanceMetaRequest>,
) -> Result<Response<LoadRebalanceMetaResponse>, Status> {
todo!()
}
async fn load_transition_tier_config(
&self,
_request: Request<LoadTransitionTierConfigRequest>,
) -> Result<Response<LoadTransitionTierConfigResponse>, Status> {
todo!()
}
}
+2
View File
@@ -1,6 +1,8 @@
mod admin;
mod config;
mod grpc;
#[allow(dead_code)]
mod peer_rest_client;
mod service;
mod storage;
+658
View File
@@ -0,0 +1,658 @@
use std::{collections::HashMap, io::Cursor, time::SystemTime};
use common::error::{Error, Result};
use ecstore::{admin_server_info::ServerProperties, store_api::StorageInfo};
use madmin::{
heal_command::BgHealState,
health::{Cpus, MemInfo, OsInfo, Partitions, ProcInfo, SysConfig, SysErrors, SysService},
metrics::{CollectMetricsOpts, MetricType, RealtimeMetrics},
net::NetInfo,
};
use protos::{
node_service_time_out_client,
proto_gen::node_service::{
BackgroundHealStatusRequest, DeleteBucketMetadataRequest, DeletePolicyRequest, DeleteServiceAccountRequest,
DeleteUserRequest, GetCpusRequest, GetMemInfoRequest, GetMetricsRequest, GetNetInfoRequest, GetOsInfoRequest,
GetPartitionsRequest, GetProcInfoRequest, GetSeLinuxInfoRequest, GetSysConfigRequest, GetSysErrorsRequest,
LoadBucketMetadataRequest, LoadGroupRequest, LoadPolicyMappingRequest, LoadPolicyRequest, LoadRebalanceMetaRequest,
LoadServiceAccountRequest, LoadTransitionTierConfigRequest, LoadUserRequest, LocalStorageInfoRequest, Mss,
ReloadPoolMetaRequest, ReloadSiteReplicationConfigRequest, ServerInfoRequest, SignalServiceRequest,
StartProfilingRequest, StopRebalanceRequest,
},
};
use rmp_serde::{Deserializer, Serializer};
use serde::{Deserialize, Serialize};
use tonic::Request;
pub const PEER_RESTSIGNAL: &str = "signal";
pub const PEER_RESTSUB_SYS: &str = "sub-sys";
pub const PEER_RESTDRY_RUN: &str = "dry-run";
struct PeerRestClient {
addr: String,
}
impl PeerRestClient {
pub fn new(url: url::Url) -> Self {
Self {
addr: format!("{}://{}:{}", url.scheme(), url.host_str().unwrap(), url.port().unwrap()),
}
}
}
impl PeerRestClient {
pub async fn local_storage_info(&self) -> Result<StorageInfo> {
let mut client = node_service_time_out_client(&self.addr)
.await
.map_err(|err| Error::msg(err.to_string()))?;
let request = Request::new(LocalStorageInfoRequest { metrics: true });
let response = client.local_storage_info(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::msg(msg));
}
return Err(Error::msg(""));
}
let data = response.storage_info;
let mut buf = Deserializer::new(Cursor::new(data));
let storage_info: StorageInfo = Deserialize::deserialize(&mut buf).unwrap();
Ok(storage_info)
}
pub async fn server_info(&self) -> Result<ServerProperties> {
let mut client = node_service_time_out_client(&self.addr)
.await
.map_err(|err| Error::msg(err.to_string()))?;
let request = Request::new(ServerInfoRequest { metrics: true });
let response = client.server_info(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::msg(msg));
}
return Err(Error::msg(""));
}
let data = response.server_properties;
let mut buf = Deserializer::new(Cursor::new(data));
let storage_properties: ServerProperties = Deserialize::deserialize(&mut buf).unwrap();
Ok(storage_properties)
}
pub async fn get_cpus(&self) -> Result<Cpus> {
let mut client = node_service_time_out_client(&self.addr)
.await
.map_err(|err| Error::msg(err.to_string()))?;
let request = Request::new(GetCpusRequest {});
let response = client.get_cpus(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::msg(msg));
}
return Err(Error::msg(""));
}
let data = response.cpus;
let mut buf = Deserializer::new(Cursor::new(data));
let cpus: Cpus = Deserialize::deserialize(&mut buf).unwrap();
Ok(cpus)
}
pub async fn get_net_info(&self) -> Result<NetInfo> {
let mut client = node_service_time_out_client(&self.addr)
.await
.map_err(|err| Error::msg(err.to_string()))?;
let request = Request::new(GetNetInfoRequest {});
let response = client.get_net_info(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::msg(msg));
}
return Err(Error::msg(""));
}
let data = response.net_info;
let mut buf = Deserializer::new(Cursor::new(data));
let net_info: NetInfo = Deserialize::deserialize(&mut buf).unwrap();
Ok(net_info)
}
pub async fn get_partitions(&self) -> Result<Partitions> {
let mut client = node_service_time_out_client(&self.addr)
.await
.map_err(|err| Error::msg(err.to_string()))?;
let request = Request::new(GetPartitionsRequest {});
let response = client.get_partitions(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::msg(msg));
}
return Err(Error::msg(""));
}
let data = response.partitions;
let mut buf = Deserializer::new(Cursor::new(data));
let partitions: Partitions = Deserialize::deserialize(&mut buf).unwrap();
Ok(partitions)
}
pub async fn get_os_info(&self) -> Result<OsInfo> {
let mut client = node_service_time_out_client(&self.addr)
.await
.map_err(|err| Error::msg(err.to_string()))?;
let request = Request::new(GetOsInfoRequest {});
let response = client.get_os_info(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::msg(msg));
}
return Err(Error::msg(""));
}
let data = response.os_info;
let mut buf = Deserializer::new(Cursor::new(data));
let os_info: OsInfo = Deserialize::deserialize(&mut buf).unwrap();
Ok(os_info)
}
pub async fn get_se_linux_info(&self) -> Result<SysService> {
let mut client = node_service_time_out_client(&self.addr)
.await
.map_err(|err| Error::msg(err.to_string()))?;
let request = Request::new(GetSeLinuxInfoRequest {});
let response = client.get_se_linux_info(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::msg(msg));
}
return Err(Error::msg(""));
}
let data = response.sys_services;
let mut buf = Deserializer::new(Cursor::new(data));
let sys_services: SysService = Deserialize::deserialize(&mut buf).unwrap();
Ok(sys_services)
}
pub async fn get_sys_config(&self) -> Result<SysConfig> {
let mut client = node_service_time_out_client(&self.addr)
.await
.map_err(|err| Error::msg(err.to_string()))?;
let request = Request::new(GetSysConfigRequest {});
let response = client.get_sys_config(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::msg(msg));
}
return Err(Error::msg(""));
}
let data = response.sys_config;
let mut buf = Deserializer::new(Cursor::new(data));
let sys_config: SysConfig = Deserialize::deserialize(&mut buf).unwrap();
Ok(sys_config)
}
pub async fn get_sys_errors(&self) -> Result<SysErrors> {
let mut client = node_service_time_out_client(&self.addr)
.await
.map_err(|err| Error::msg(err.to_string()))?;
let request = Request::new(GetSysErrorsRequest {});
let response = client.get_sys_errors(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::msg(msg));
}
return Err(Error::msg(""));
}
let data = response.sys_errors;
let mut buf = Deserializer::new(Cursor::new(data));
let sys_errors: SysErrors = Deserialize::deserialize(&mut buf).unwrap();
Ok(sys_errors)
}
pub async fn get_mem_info(&self) -> Result<MemInfo> {
let mut client = node_service_time_out_client(&self.addr)
.await
.map_err(|err| Error::msg(err.to_string()))?;
let request = Request::new(GetMemInfoRequest {});
let response = client.get_mem_info(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::msg(msg));
}
return Err(Error::msg(""));
}
let data = response.mem_info;
let mut buf = Deserializer::new(Cursor::new(data));
let mem_info: MemInfo = Deserialize::deserialize(&mut buf).unwrap();
Ok(mem_info)
}
pub async fn get_metrics(&self, t: MetricType, opts: &CollectMetricsOpts) -> Result<RealtimeMetrics> {
let mut client = node_service_time_out_client(&self.addr)
.await
.map_err(|err| Error::msg(err.to_string()))?;
let mut buf = Vec::new();
opts.serialize(&mut Serializer::new(&mut buf))?;
let request = Request::new(GetMetricsRequest {
metric_type: t,
opts: buf,
});
let response = client.get_metrics(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::msg(msg));
}
return Err(Error::msg(""));
}
let data = response.realtime_metrics;
let mut buf = Deserializer::new(Cursor::new(data));
let realtime_metrics: RealtimeMetrics = Deserialize::deserialize(&mut buf).unwrap();
Ok(realtime_metrics)
}
pub async fn get_proc_info(&self) -> Result<ProcInfo> {
let mut client = node_service_time_out_client(&self.addr)
.await
.map_err(|err| Error::msg(err.to_string()))?;
let request = Request::new(GetProcInfoRequest {});
let response = client.get_proc_info(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::msg(msg));
}
return Err(Error::msg(""));
}
let data = response.proc_info;
let mut buf = Deserializer::new(Cursor::new(data));
let proc_info: ProcInfo = Deserialize::deserialize(&mut buf).unwrap();
Ok(proc_info)
}
pub async fn start_profiling(&self, profiler: &str) -> Result<()> {
let mut client = node_service_time_out_client(&self.addr)
.await
.map_err(|err| Error::msg(err.to_string()))?;
let request = Request::new(StartProfilingRequest {
profiler: profiler.to_string(),
});
let response = client.start_profiling(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::msg(msg));
}
return Err(Error::msg(""));
}
Ok(())
}
pub async fn download_profile_data(&self) -> Result<()> {
todo!()
}
pub async fn get_bucket_stats(&self) -> Result<()> {
todo!()
}
pub async fn get_sr_metrics(&self) -> Result<()> {
todo!()
}
pub async fn get_all_bucket_stats(&self) -> Result<()> {
todo!()
}
pub async fn load_bucket_metadata(&self, bucket: &str) -> Result<()> {
let mut client = node_service_time_out_client(&self.addr)
.await
.map_err(|err| Error::msg(err.to_string()))?;
let request = Request::new(LoadBucketMetadataRequest {
bucket: bucket.to_string(),
});
let response = client.load_bucket_metadata(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::msg(msg));
}
return Err(Error::msg(""));
}
Ok(())
}
pub async fn delete_bucket_metadata(&self, bucket: &str) -> Result<()> {
let mut client = node_service_time_out_client(&self.addr)
.await
.map_err(|err| Error::msg(err.to_string()))?;
let request = Request::new(DeleteBucketMetadataRequest {
bucket: bucket.to_string(),
});
let response = client.delete_bucket_metadata(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::msg(msg));
}
return Err(Error::msg(""));
}
Ok(())
}
pub async fn delete_policy(&self, policy: &str) -> Result<()> {
let mut client = node_service_time_out_client(&self.addr)
.await
.map_err(|err| Error::msg(err.to_string()))?;
let request = Request::new(DeletePolicyRequest {
policy_name: policy.to_string(),
});
let response = client.delete_policy(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::msg(msg));
}
return Err(Error::msg(""));
}
Ok(())
}
pub async fn load_policy(&self, policy: &str) -> Result<()> {
let mut client = node_service_time_out_client(&self.addr)
.await
.map_err(|err| Error::msg(err.to_string()))?;
let request = Request::new(LoadPolicyRequest {
policy_name: policy.to_string(),
});
let response = client.load_policy(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::msg(msg));
}
return Err(Error::msg(""));
}
Ok(())
}
pub async fn load_policy_mapping(&self, user_or_group: &str, user_type: u64, is_group: bool) -> Result<()> {
let mut client = node_service_time_out_client(&self.addr)
.await
.map_err(|err| Error::msg(err.to_string()))?;
let request = Request::new(LoadPolicyMappingRequest {
user_or_group: user_or_group.to_string(),
user_type,
is_group,
});
let response = client.load_policy_mapping(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::msg(msg));
}
return Err(Error::msg(""));
}
Ok(())
}
pub async fn delete_user(&self, access_key: &str) -> Result<()> {
let mut client = node_service_time_out_client(&self.addr)
.await
.map_err(|err| Error::msg(err.to_string()))?;
let request = Request::new(DeleteUserRequest {
access_key: access_key.to_string(),
});
let response = client.delete_user(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::msg(msg));
}
return Err(Error::msg(""));
}
Ok(())
}
pub async fn delete_service_account(&self, access_key: &str) -> Result<()> {
let mut client = node_service_time_out_client(&self.addr)
.await
.map_err(|err| Error::msg(err.to_string()))?;
let request = Request::new(DeleteServiceAccountRequest {
access_key: access_key.to_string(),
});
let response = client.delete_service_account(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::msg(msg));
}
return Err(Error::msg(""));
}
Ok(())
}
pub async fn load_user(&self, access_key: &str, temp: bool) -> Result<()> {
let mut client = node_service_time_out_client(&self.addr)
.await
.map_err(|err| Error::msg(err.to_string()))?;
let request = Request::new(LoadUserRequest {
access_key: access_key.to_string(),
temp,
});
let response = client.load_user(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::msg(msg));
}
return Err(Error::msg(""));
}
Ok(())
}
pub async fn load_service_account(&self, access_key: &str) -> Result<()> {
let mut client = node_service_time_out_client(&self.addr)
.await
.map_err(|err| Error::msg(err.to_string()))?;
let request = Request::new(LoadServiceAccountRequest {
access_key: access_key.to_string(),
});
let response = client.load_service_account(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::msg(msg));
}
return Err(Error::msg(""));
}
Ok(())
}
pub async fn load_group(&self, group: &str) -> Result<()> {
let mut client = node_service_time_out_client(&self.addr)
.await
.map_err(|err| Error::msg(err.to_string()))?;
let request = Request::new(LoadGroupRequest {
group: group.to_string(),
});
let response = client.load_group(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::msg(msg));
}
return Err(Error::msg(""));
}
Ok(())
}
pub async fn reload_site_replication_config(&self) -> Result<()> {
let mut client = node_service_time_out_client(&self.addr)
.await
.map_err(|err| Error::msg(err.to_string()))?;
let request = Request::new(ReloadSiteReplicationConfigRequest {});
let response = client.reload_site_replication_config(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::msg(msg));
}
return Err(Error::msg(""));
}
Ok(())
}
pub async fn signal_service(&self, sig: u64, sub_sys: &str, dry_run: bool, _exec_at: SystemTime) -> Result<()> {
let mut client = node_service_time_out_client(&self.addr)
.await
.map_err(|err| Error::msg(err.to_string()))?;
let mut vars = HashMap::new();
vars.insert(PEER_RESTSIGNAL.to_string(), sig.to_string());
vars.insert(PEER_RESTSUB_SYS.to_string(), sub_sys.to_string());
vars.insert(PEER_RESTDRY_RUN.to_string(), dry_run.to_string());
let request = Request::new(SignalServiceRequest {
vars: Some(Mss { value: vars }),
});
let response = client.signal_service(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::msg(msg));
}
return Err(Error::msg(""));
}
Ok(())
}
pub async fn background_heal_status(&self) -> Result<BgHealState> {
let mut client = node_service_time_out_client(&self.addr)
.await
.map_err(|err| Error::msg(err.to_string()))?;
let request = Request::new(BackgroundHealStatusRequest {});
let response = client.background_heal_status(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::msg(msg));
}
return Err(Error::msg(""));
}
let data = response.bg_heal_state;
let mut buf = Deserializer::new(Cursor::new(data));
let bg_heal_state: BgHealState = Deserialize::deserialize(&mut buf).unwrap();
Ok(bg_heal_state)
}
pub async fn get_metacache_listing(&self) -> Result<()> {
let mut _client = node_service_time_out_client(&self.addr)
.await
.map_err(|err| Error::msg(err.to_string()))?;
todo!()
}
pub async fn update_metacache_listing(&self) -> Result<()> {
let mut _client = node_service_time_out_client(&self.addr)
.await
.map_err(|err| Error::msg(err.to_string()))?;
todo!()
}
pub async fn reload_pool_meta(&self) -> Result<()> {
let mut client = node_service_time_out_client(&self.addr)
.await
.map_err(|err| Error::msg(err.to_string()))?;
let request = Request::new(ReloadPoolMetaRequest {});
let response = client.reload_pool_meta(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::msg(msg));
}
return Err(Error::msg(""));
}
Ok(())
}
pub async fn stop_rebalance(&self) -> Result<()> {
let mut client = node_service_time_out_client(&self.addr)
.await
.map_err(|err| Error::msg(err.to_string()))?;
let request = Request::new(StopRebalanceRequest {});
let response = client.stop_rebalance(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::msg(msg));
}
return Err(Error::msg(""));
}
Ok(())
}
pub async fn load_rebalance_meta(&self, start_rebalance: bool) -> Result<()> {
let mut client = node_service_time_out_client(&self.addr)
.await
.map_err(|err| Error::msg(err.to_string()))?;
let request = Request::new(LoadRebalanceMetaRequest { start_rebalance });
let response = client.load_rebalance_meta(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::msg(msg));
}
return Err(Error::msg(""));
}
Ok(())
}
pub async fn load_transition_tier_config(&self) -> Result<()> {
let mut client = node_service_time_out_client(&self.addr)
.await
.map_err(|err| Error::msg(err.to_string()))?;
let request = Request::new(LoadTransitionTierConfigRequest {});
let response = client.load_transition_tier_config(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::msg(msg));
}
return Err(Error::msg(""));
}
Ok(())
}
}