support some peer rest api

Signed-off-by: mujunxiang <1948535941@qq.com>
This commit is contained in:
mujunxiang
2024-11-27 14:21:26 +08:00
parent 5cc138e23d
commit 1b9ae3ccb3
16 changed files with 1482 additions and 68 deletions
Generated
+1
View File
@@ -1454,6 +1454,7 @@ dependencies = [
name = "madmin"
version = "0.0.1"
dependencies = [
"ecstore",
"psutil",
"serde",
]
@@ -1036,10 +1036,8 @@ pub struct LoadRebalanceMetaResponse {
pub error_info: ::core::option::Option<::prost::alloc::string::String>,
}
#[derive(Clone, Copy, PartialEq, ::prost::Message)]
pub struct LoadTransitionTierConfigRequest {
#[prost(bool, tag = "1")]
pub start_rebalance: bool,
}
pub struct LoadTransitionTierConfigRequest {}
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct LoadTransitionTierConfigResponse {
#[prost(bool, tag = "1")]
+1 -3
View File
@@ -725,9 +725,7 @@ message LoadRebalanceMetaResponse {
optional string error_info = 2;
}
message LoadTransitionTierConfigRequest {
bool start_rebalance = 1;
}
message LoadTransitionTierConfigRequest {}
message LoadTransitionTierConfigResponse {
bool success = 1;
+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 {
+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()
}
+1 -1
View File
@@ -9,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;
+107 -13
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,
@@ -157,15 +225,11 @@ impl ECStore {
}
}
fn get_total_usable_capacity(disks: &Vec<StorageDisk>, info: &StorageInfo) -> usize {
let mut capacity = 0;
for disk in disks.iter() {
if disk.pool_index < 0 || info.backend.standard_sc_data.len() <= disk.pool_index as usize {
continue;
}
if (disk.disk_index as usize) < info.backend.standard_sc_data[disk.pool_index as usize] {
capacity += disk.total_space as usize;
}
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;
// }
}
capacity
}
@@ -182,3 +246,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(())
}
+34 -7
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> {
@@ -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(
+2 -1
View File
@@ -7,5 +7,6 @@ rust-version.workspace = true
version.workspace = true
[dependencies]
ecstore.workspace = true
psutil = "3.3.0"
serde.workspace = true
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)
}
+126
View File
@@ -1,3 +1,5 @@
use std::collections::HashMap;
use serde::{Deserialize, Serialize};
#[derive(Debug, Default, Serialize, Deserialize)]
@@ -49,3 +51,127 @@ 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()
}
+2
View File
@@ -1,2 +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()
}
+1
View File
@@ -2,6 +2,7 @@ use serde::{Deserialize, Serialize};
use crate::health::NodeCommon;
#[cfg(target_os = "linux")]
pub mod net_linux;
#[derive(Debug, Default, Serialize, Deserialize)]
+442 -29
View File
@@ -1,7 +1,14 @@
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,
@@ -16,7 +23,15 @@ use ecstore::{
use futures::{Stream, StreamExt};
use lock::{lock_args::LockArgs, Locker, GLOBAL_LOCAL_SERVER};
use madmin::health::get_cpus;
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, *},
@@ -1569,42 +1584,168 @@ impl Node for NodeService {
}
async fn get_net_info(&self, _request: Request<GetNetInfoRequest>) -> Result<Response<GetNetInfoResponse>, Status> {
todo!()
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> {
todo!()
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> {
todo!()
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> {
todo!()
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> {
todo!()
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> {
todo!()
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> {
todo!()
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> {
todo!()
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> {
todo!()
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(
@@ -1644,56 +1785,257 @@ impl Node for NodeService {
async fn load_bucket_metadata(
&self,
_request: Request<LoadBucketMetadataRequest>,
request: Request<LoadBucketMetadataRequest>,
) -> Result<Response<LoadBucketMetadataResponse>, Status> {
todo!()
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>,
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 delete_policy(&self, _request: Request<DeletePolicyRequest>) -> Result<Response<DeletePolicyResponse>, Status> {
todo!()
}
async fn load_policy(&self, _request: Request<LoadPolicyRequest>) -> Result<Response<LoadPolicyResponse>, Status> {
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>,
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> {
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>,
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> {
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>,
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> {
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!()
}
@@ -1701,10 +2043,26 @@ impl Node for NodeService {
&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> {
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!()
}
@@ -1712,7 +2070,28 @@ impl Node for NodeService {
&self,
_request: Request<BackgroundHealStatusRequest>,
) -> Result<Response<BackgroundHealStatusResponse>, Status> {
todo!()
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(
@@ -1733,10 +2112,44 @@ impl Node for NodeService {
&self,
_request: Request<ReloadPoolMetaRequest>,
) -> Result<Response<ReloadPoolMetaResponse>, Status> {
todo!()
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!()
}
+574 -5
View File
@@ -1,16 +1,33 @@
use std::io::Cursor;
use std::{collections::HashMap, io::Cursor, time::SystemTime};
use common::error::{Error, Result};
use ecstore::{admin_server_info::ServerProperties, store_api::StorageInfo};
use madmin::health::Cpus;
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::{GetCpusRequest, LocalStorageInfoRequest, ServerInfoRequest},
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;
use serde::Deserialize;
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,
}
@@ -86,4 +103,556 @@ impl PeerRestClient {
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, 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(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(())
}
}