diff --git a/Cargo.lock b/Cargo.lock index 0d17bba47..0422aed7d 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1454,6 +1454,7 @@ dependencies = [ name = "madmin" version = "0.0.1" dependencies = [ + "ecstore", "psutil", "serde", ] diff --git a/common/protos/src/generated/proto_gen/node_service.rs b/common/protos/src/generated/proto_gen/node_service.rs index 29f002936..9dc50d8cc 100644 --- a/common/protos/src/generated/proto_gen/node_service.rs +++ b/common/protos/src/generated/proto_gen/node_service.rs @@ -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")] diff --git a/common/protos/src/node.proto b/common/protos/src/node.proto index 5cd4015a2..e3c0ad667 100644 --- a/common/protos/src/node.proto +++ b/common/protos/src/node.proto @@ -725,9 +725,7 @@ message LoadRebalanceMetaResponse { optional string error_info = 2; } -message LoadTransitionTierConfigRequest { - bool start_rebalance = 1; -} +message LoadTransitionTierConfigRequest {} message LoadTransitionTierConfigResponse { bool success = 1; diff --git a/ecstore/src/heal/background_heal_ops.rs b/ecstore/src/heal/background_heal_ops.rs index c60287b5a..0c64482a1 100644 --- a/ecstore/src/heal/background_heal_ops.rs +++ b/ecstore/src/heal/background_heal_ops.rs @@ -65,7 +65,7 @@ async fn init_background_healing() { .await; } -async fn get_local_disks_to_heal() -> Vec { +pub async fn get_local_disks_to_heal() -> Vec { let mut disks_to_heal = Vec::new(); for (_, disk) in GLOBAL_LOCAL_DISK_MAP.read().await.iter() { if let Some(disk) = disk { diff --git a/ecstore/src/heal/heal_ops.rs b/ecstore/src/heal/heal_ops.rs index c869e4ebb..a271f6c36 100644 --- a/ecstore/src/heal/heal_ops.rs +++ b/ecstore/src/heal/heal_ops.rs @@ -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() } diff --git a/ecstore/src/lib.rs b/ecstore/src/lib.rs index 44d8cd0e1..9c5bbaa3a 100644 --- a/ecstore/src/lib.rs +++ b/ecstore/src/lib.rs @@ -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; diff --git a/ecstore/src/pools.rs b/ecstore/src/pools.rs index 53a22392b..da57427b7 100644 --- a/ecstore/src/pools.rs +++ b/ecstore/src/pools.rs @@ -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, } -#[derive(Debug, Clone)] +#[derive(Debug, Clone, Default, Serialize, Deserialize)] pub struct PoolMeta { + pub version: u16, pub pools: Vec, 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::() { + 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::(POOL_META_FORMAT).unwrap(); + data.write_u16::(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, pub start_size: usize, @@ -157,15 +225,11 @@ impl ECStore { } } -fn get_total_usable_capacity(disks: &Vec, 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, _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, info: &StorageInfo) } capacity } + +#[test] +fn test_pool_meta() -> Result<()> { + let meta = PoolMeta::new(vec![]); + let mut data = Vec::new(); + data.write_u16::(POOL_META_FORMAT).unwrap(); + data.write_u16::(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(()) +} diff --git a/ecstore/src/store.rs b/ecstore/src/store.rs index 927bcdc7a..0ac2cca62 100644 --- a/ecstore/src/store.rs +++ b/ecstore/src/store.rs @@ -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, @@ -76,10 +76,27 @@ pub struct ECStore { pub pools: Vec>, pub peer_sys: S3PeerSys, // pub local_disks: Vec, - pub pool_meta: PoolMeta, + pub pool_meta: std_RwLock, pub decommission_cancelers: Vec>, } +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 { @@ -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 { @@ -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, opts: &ObjectOptions) -> Vec { 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( diff --git a/madmin/Cargo.toml b/madmin/Cargo.toml index 05e052f1f..d0b8cd6f5 100644 --- a/madmin/Cargo.toml +++ b/madmin/Cargo.toml @@ -7,5 +7,6 @@ rust-version.workspace = true version.workspace = true [dependencies] +ecstore.workspace = true psutil = "3.3.0" -serde.workspace = true \ No newline at end of file +serde.workspace = true diff --git a/madmin/src/heal_command.rs b/madmin/src/heal_command.rs new file mode 100644 index 000000000..ad87352af --- /dev/null +++ b/madmin/src/heal_command.rs @@ -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, +} + +#[derive(Debug, Default, Serialize, Deserialize)] +pub struct BgHealState { + offline_endpoints: Vec, + scanned_items_count: u64, + heal_disks: Vec, + sets: Vec, + mrf: HashMap, + scparity: HashMap, +} + +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) +} diff --git a/madmin/src/health.rs b/madmin/src/health.rs index f6d778a70..be9f5cb4c 100644 --- a/madmin/src/health.rs +++ b/madmin/src/health.rs @@ -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, +} + +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, + cmd_line: String, + num_connections: usize, + create_time: u64, + cwd: String, + exec_path: String, + gids: Vec, + // 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, + 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, +} + +pub fn get_sys_services(_add: &str) -> SysServices { + SysServices::default() +} + +#[derive(Debug, Default, Serialize, Deserialize)] +pub struct SysConfig { + node_common: NodeCommon, + config: HashMap, +} + +pub fn get_sys_config(_addr: &str) -> SysConfig { + SysConfig::default() +} + +#[derive(Debug, Default, Serialize, Deserialize)] +pub struct SysErrors { + node_common: NodeCommon, + errors: Vec, +} + +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() +} diff --git a/madmin/src/lib.rs b/madmin/src/lib.rs index d36b2f8d7..11da4062f 100644 --- a/madmin/src/lib.rs +++ b/madmin/src/lib.rs @@ -1,2 +1,4 @@ +pub mod heal_command; pub mod health; +pub mod metrics; pub mod net; diff --git a/madmin/src/metrics.rs b/madmin/src/metrics.rs new file mode 100644 index 000000000..15abc0bfc --- /dev/null +++ b/madmin/src/metrics.rs @@ -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, + last_minute: HashMap, + io_stats: DiskIOStats, +} + +#[derive(Debug, Default, Serialize, Deserialize)] +pub struct Metrics {} + +#[derive(Debug, Default, Serialize, Deserialize)] +pub struct RealtimeMetrics { + errors: Vec, + hosts: Vec, + aggregated: Metrics, + by_host: HashMap, + by_disk: HashMap, + finally: bool, +} + +#[derive(Debug, Default, Serialize, Deserialize)] +pub struct CollectMetricsOpts { + hosts: HashSet, + disks: HashSet, + job_id: String, + dep_id: String, +} + +pub type MetricType = u64; + +pub fn collect_local_metrics(_types: MetricType, _opts: &CollectMetricsOpts) -> RealtimeMetrics { + RealtimeMetrics::default() +} diff --git a/madmin/src/net/mod.rs b/madmin/src/net/mod.rs index 3f22342d4..6b06e5f45 100644 --- a/madmin/src/net/mod.rs +++ b/madmin/src/net/mod.rs @@ -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)] diff --git a/rustfs/src/grpc.rs b/rustfs/src/grpc.rs index c05b3f63d..da9a744e9 100644 --- a/rustfs/src/grpc.rs +++ b/rustfs/src/grpc.rs @@ -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) -> Result, 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) -> Result, 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) -> Result, 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, ) -> Result, 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) -> Result, 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) -> Result, 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) -> Result, 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) -> Result, Status> { - todo!() + async fn get_metrics(&self, request: Request) -> Result, 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) -> Result, 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, + request: Request, ) -> Result, 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, + request: Request, ) -> Result, 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) -> Result, 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) -> Result, Status> { - todo!() - } - - async fn load_policy(&self, _request: Request) -> Result, Status> { + async fn load_policy(&self, request: Request) -> Result, 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, + request: Request, ) -> Result, 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) -> Result, Status> { + async fn delete_user(&self, request: Request) -> Result, 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, + request: Request, ) -> Result, 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) -> Result, Status> { + async fn load_user(&self, request: Request) -> Result, 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, + request: Request, ) -> Result, 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) -> Result, Status> { + async fn load_group(&self, request: Request) -> Result, 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, ) -> Result, 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) -> Result, Status> { + async fn signal_service(&self, request: Request) -> Result, 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, ) -> Result, 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, ) -> Result, 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) -> Result, 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!() } diff --git a/rustfs/src/peer_rest_client.rs b/rustfs/src/peer_rest_client.rs index 24cee0050..d24f17a9d 100644 --- a/rustfs/src/peer_rest_client.rs +++ b/rustfs/src/peer_rest_client.rs @@ -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 { + 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 { + 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 { + 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 { + 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 { + 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 { + 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 { + 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 { + 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 { + 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 { + 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(()) + } }