From 6ddfc809438479b42eb7785ae11d3c6f2ea92a86 Mon Sep 17 00:00:00 2001 From: mujunxiang <1948535941@qq.com> Date: Sat, 23 Nov 2024 14:18:35 +0800 Subject: [PATCH] heal admin api(2) Signed-off-by: mujunxiang <1948535941@qq.com> --- ecstore/src/disk/error.rs | 15 +- ecstore/src/disk/local.rs | 2 + ecstore/src/disk/mod.rs | 1 + ecstore/src/global.rs | 2 +- ecstore/src/heal/background_heal_ops.rs | 32 ++-- ecstore/src/heal/data_scanner.rs | 6 - ecstore/src/heal/heal_ops.rs | 227 ++++++++++++------------ ecstore/src/peer.rs | 15 +- ecstore/src/sets.rs | 2 +- ecstore/src/store.rs | 4 +- ecstore/src/store_api.rs | 5 +- rustfs/src/main.rs | 6 +- 12 files changed, 167 insertions(+), 150 deletions(-) diff --git a/ecstore/src/disk/error.rs b/ecstore/src/disk/error.rs index 0bce48455..0902e5d05 100644 --- a/ecstore/src/disk/error.rs +++ b/ecstore/src/disk/error.rs @@ -419,15 +419,18 @@ pub fn convert_access_error(e: io::Error, per_err: DiskError) -> Error { } pub fn is_all_not_found(errs: &[Option]) -> bool { - for err in errs.iter().flatten() { - if let Some(err) = err.downcast_ref::() { - match err { - DiskError::FileNotFound | DiskError::VolumeNotFound | &DiskError::FileVersionNotFound => { - continue; + for err in errs.iter() { + if let Some(err) = err { + if let Some(err) = err.downcast_ref::() { + match err { + DiskError::FileNotFound | DiskError::VolumeNotFound | &DiskError::FileVersionNotFound => { + continue; + } + _ => return false, } - _ => return false, } } + return false; } !errs.is_empty() diff --git a/ecstore/src/disk/local.rs b/ecstore/src/disk/local.rs index 46dee48a6..e1d82a938 100644 --- a/ecstore/src/disk/local.rs +++ b/ecstore/src/disk/local.rs @@ -1172,6 +1172,7 @@ impl DiskAPI for LocalDisk { check_path_length(file_path.to_string_lossy().to_string().as_str())?; // TODO: writeAllDirect io.copy + info!("file_path: {:?}", file_path); if let Some(parent) = file_path.parent() { os::make_dir_all(parent, &volume_dir).await?; } @@ -1924,6 +1925,7 @@ impl DiskAPI for LocalDisk { } async fn delete_volume(&self, volume: &str) -> Result<()> { + info!("delete_volume, volume: {}", volume); let p = self.get_bucket_path(volume)?; // TODO: 不能用递归删除,如果目录下面有文件,返回errVolumeNotEmpty diff --git a/ecstore/src/disk/mod.rs b/ecstore/src/disk/mod.rs index 10ceeafb7..0703e59dc 100644 --- a/ecstore/src/disk/mod.rs +++ b/ecstore/src/disk/mod.rs @@ -330,6 +330,7 @@ impl DiskAPI for Disk { } async fn delete_volume(&self, volume: &str) -> Result<()> { + info!("delete_volume, volume: {}", volume); match self { Disk::Local(local_disk) => local_disk.delete_volume(volume).await, Disk::Remote(remote_disk) => remote_disk.delete_volume(volume).await, diff --git a/ecstore/src/global.rs b/ecstore/src/global.rs index fdb7f8444..5bb5f0153 100644 --- a/ecstore/src/global.rs +++ b/ecstore/src/global.rs @@ -28,7 +28,7 @@ lazy_static! { pub static ref GLOBAL_LOCAL_DISK_SET_DRIVES: Arc> = Arc::new(RwLock::new(Vec::new())); pub static ref GLOBAL_Endpoints: OnceLock = OnceLock::new(); pub static ref GLOBAL_RootDiskThreshold: RwLock = RwLock::new(0); - pub static ref GLOBAL_BackgroundHealRoutine: Arc> = HealRoutine::new(); + pub static ref GLOBAL_BackgroundHealRoutine: Arc = HealRoutine::new(); pub static ref GLOBAL_BackgroundHealState: Arc> = AllHealState::new(false); static ref globalDeploymentIDPtr: RwLock = RwLock::new(Uuid::nil()); } diff --git a/ecstore/src/heal/background_heal_ops.rs b/ecstore/src/heal/background_heal_ops.rs index cdd33f57e..14f44eb15 100644 --- a/ecstore/src/heal/background_heal_ops.rs +++ b/ecstore/src/heal/background_heal_ops.rs @@ -51,11 +51,11 @@ pub async fn init_auto_heal() { } async fn init_background_healing() { - let bg_seq = Arc::new(RwLock::new(new_bg_heal_sequence())); - for _ in 0..GLOBAL_BackgroundHealRoutine.read().await.workers { + let bg_seq = Arc::new(new_bg_heal_sequence()); + for _ in 0..GLOBAL_BackgroundHealRoutine.workers { let bg_seq_clone = bg_seq.clone(); tokio::spawn(async { - GLOBAL_BackgroundHealRoutine.write().await.add_worker(bg_seq_clone).await; + GLOBAL_BackgroundHealRoutine.add_worker(bg_seq_clone).await; }); } let _ = GLOBAL_BackgroundHealState @@ -318,6 +318,7 @@ impl HealTask { } } +#[derive(Debug)] pub struct HealResult { pub result: HealResultItem, pub err: Option, @@ -325,12 +326,12 @@ pub struct HealResult { pub struct HealRoutine { pub tasks_tx: Sender, - tasks_rx: Receiver, + tasks_rx: RwLock>, workers: usize, } impl HealRoutine { - pub fn new() -> Arc> { + pub fn new() -> Arc { let mut workers = num_cpus::get() / 2; if let Ok(env_heal_workers) = env::var("_RUSTFS_HEAL_WORKERS") { if let Ok(num_healers) = env_heal_workers.parse::() { @@ -343,19 +344,20 @@ impl HealRoutine { } let (tx, rx) = mpsc::channel(100); - Arc::new(RwLock::new(Self { + Arc::new(Self { tasks_tx: tx, - tasks_rx: rx, + tasks_rx: RwLock::new(rx), workers, - })) + }) } - pub async fn add_worker(&mut self, bgseq: Arc>) { + pub async fn add_worker(&self, bgseq: Arc) { loop { let mut d_res = HealResultItem::default(); let d_err: Option; - match self.tasks_rx.recv().await { + match self.tasks_rx.write().await.recv().await { Some(task) => { + info!("got task: {:?}", task); if task.bucket == NOP_HEAL { d_err = Some(Error::from_string("skip file")); } else if task.bucket == SLASH_SEPARATOR { @@ -389,6 +391,7 @@ impl HealRoutine { } } } + info!("task finished, task: {:?}", task); if let Some(resp_tx) = task.resp_tx { let _ = resp_tx .send(HealResult { @@ -400,13 +403,16 @@ impl HealRoutine { // when respCh is not set caller is not waiting but we // update the relevant metrics for them if d_err.is_none() { - bgseq.write().await.count_healed(d_res.heal_item_type); + bgseq.count_healed(d_res.heal_item_type).await; } else { - bgseq.write().await.count_failed(d_res.heal_item_type); + bgseq.count_failed(d_res.heal_item_type).await; } } } - None => return, + None => { + info!("add_worker, tasks_rx was closed, return"); + return; + }, } } } diff --git a/ecstore/src/heal/data_scanner.rs b/ecstore/src/heal/data_scanner.rs index bb25cec01..3f1656969 100644 --- a/ecstore/src/heal/data_scanner.rs +++ b/ecstore/src/heal/data_scanner.rs @@ -748,8 +748,6 @@ impl FolderScanner { if bucket != resolver.bucket { bg_seq .clone() - .write() - .await .queue_heal_task( HealSource { bucket: bucket.clone(), @@ -818,8 +816,6 @@ impl FolderScanner { Ok(fiv) => fiv, Err(_) => { if let Err(err) = bg_seq_partial - .write() - .await .queue_heal_task( HealSource { bucket: bucket_partial.clone(), @@ -849,8 +845,6 @@ impl FolderScanner { let (mut success_versions, mut fail_versions) = (0, 0); for ver in fiv.versions.iter() { match bg_seq_partial - .write() - .await .queue_heal_task( HealSource { bucket: bucket_partial.clone(), diff --git a/ecstore/src/heal/heal_ops.rs b/ecstore/src/heal/heal_ops.rs index 2913e28c9..d5a3555e7 100644 --- a/ecstore/src/heal/heal_ops.rs +++ b/ecstore/src/heal/heal_ops.rs @@ -7,7 +7,7 @@ use super::{ HEAL_ITEM_BUCKET_METADATA, }, }; -use crate::heal::heal_commands::{HEAL_DEEP_SCAN, HEAL_ITEM_BUCKET, HEAL_ITEM_OBJECT}; +use crate::heal::heal_commands::{HEAL_ITEM_BUCKET, HEAL_ITEM_OBJECT}; use crate::store_api::StorageAPI; use crate::{ config::common::CONFIG_PREFIX, @@ -27,6 +27,7 @@ use crate::{ new_object_layer_fn, utils::path::has_profix, }; +use futures::join; use lazy_static::lazy_static; use std::{ collections::HashMap, @@ -39,13 +40,14 @@ use std::{ use tokio::{ select, spawn, sync::{ + watch::{self, Receiver as W_Receiver, Sender as W_Sender}, broadcast::{self, Receiver, Sender}, mpsc::{self, Receiver as M_Receiver, Sender as M_Sender}, RwLock, }, time::{interval, sleep}, }; -use tracing::info; +use tracing::{error, info}; use uuid::Uuid; type HealStatusSummary = String; @@ -89,7 +91,7 @@ pub struct HealSource { pub opts: Option, } -#[derive(Clone, Debug)] +#[derive(Debug)] pub struct HealSequence { pub bucket: String, pub object: String, @@ -102,16 +104,16 @@ pub struct HealSequence { pub setting: HealOpts, pub current_status: Arc>, pub last_sent_result_index: usize, - pub scanned_items_map: ItemsMap, - pub healed_items_map: ItemsMap, - pub heal_failed_items_map: ItemsMap, - pub last_heal_activity: u64, + pub scanned_items_map: RwLock, + pub healed_items_map: RwLock, + pub heal_failed_items_map: RwLock, + pub last_heal_activity: RwLock, traverse_and_heal_done_tx: Arc>>>, traverse_and_heal_done_rx: Arc>>>, - tx: Arc>>, - rx: Arc>>, + tx: W_Sender, + rx: W_Receiver, } pub fn new_bg_heal_sequence() -> HealSequence { @@ -131,9 +133,9 @@ pub fn new_bg_heal_sequence() -> HealSequence { ..Default::default() })), report_progress: false, - scanned_items_map: HashMap::new(), - healed_items_map: HashMap::new(), - heal_failed_items_map: HashMap::new(), + scanned_items_map: HashMap::new().into(), + healed_items_map: HashMap::new().into(), + heal_failed_items_map: HashMap::new().into(), ..Default::default() } } @@ -157,9 +159,9 @@ pub fn new_heal_sequence(bucket: &str, obj_prefix: &str, client_addr: &str, hs: })), traverse_and_heal_done_tx: Arc::new(RwLock::new(tx)), traverse_and_heal_done_rx: Arc::new(RwLock::new(rx)), - scanned_items_map: HashMap::new(), - healed_items_map: HashMap::new(), - heal_failed_items_map: HashMap::new(), + scanned_items_map: HashMap::new().into(), + healed_items_map: HashMap::new().into(), + heal_failed_items_map: HashMap::new().into(), ..Default::default() } } @@ -167,7 +169,7 @@ pub fn new_heal_sequence(bucket: &str, obj_prefix: &str, client_addr: &str, hs: impl Default for HealSequence { fn default() -> Self { let (h_tx, h_rx) = mpsc::channel(1); - let (tx, rx) = broadcast::channel(1); + let (tx, rx) = watch::channel(false); Self { bucket: Default::default(), object: Default::default(), @@ -183,11 +185,11 @@ impl Default for HealSequence { scanned_items_map: Default::default(), healed_items_map: Default::default(), heal_failed_items_map: Default::default(), - last_heal_activity: Default::default(), + last_heal_activity: RwLock::new(SystemTime::now()), traverse_and_heal_done_tx: Arc::new(RwLock::new(h_tx)), traverse_and_heal_done_rx: Arc::new(RwLock::new(h_rx)), - tx: Arc::new(RwLock::new(tx)), - rx: Arc::new(RwLock::new(rx)), + tx, + rx, } } } @@ -217,47 +219,42 @@ impl HealSequence { impl HealSequence { pub fn get_scanned_items_count(&self) -> usize { self.scanned_items_map.values().sum() + async fn _get_scanned_items_count(&self) -> usize { + self.scanned_items_map.read().await.values().sum() } pub fn _get_scanned_items_map(&self) -> ItemsMap { self.scanned_items_map.clone() + async fn _get_scanned_items_map(&self) -> ItemsMap { + self.scanned_items_map.read().await.clone() } - pub fn _get_healed_items_map(&self) -> ItemsMap { - self.healed_items_map.clone() + async fn _get_healed_items_map(&self) -> ItemsMap { + self.healed_items_map.read().await.clone() } - pub fn _get_heal_failed_items_map(&self) -> ItemsMap { - self.heal_failed_items_map.clone() + async fn _get_heal_failed_items_map(&self) -> ItemsMap { + self.heal_failed_items_map.read().await.clone() } - pub fn count_failed(&mut self, heal_type: HealItemType) { - *self.heal_failed_items_map.entry(heal_type).or_insert(0) += 1; - self.last_heal_activity = SystemTime::now() - .duration_since(UNIX_EPOCH) - .expect("Time went backwards") - .as_secs(); + pub async fn count_failed(&self, heal_type: HealItemType) { + *self.heal_failed_items_map.write().await.entry(heal_type).or_insert(0) += 1; + *self.last_heal_activity.write().await = SystemTime::now(); } - pub fn count_scanned(&mut self, heal_type: HealItemType) { - *self.scanned_items_map.entry(heal_type).or_insert(0) += 1; - self.last_heal_activity = SystemTime::now() - .duration_since(UNIX_EPOCH) - .expect("Time went backwards") - .as_secs(); + pub async fn count_scanned(&self, heal_type: HealItemType) { + *self.scanned_items_map.write().await.entry(heal_type).or_insert(0) += 1; + *self.last_heal_activity.write().await = SystemTime::now(); } - pub fn count_healed(&mut self, heal_type: HealItemType) { - *self.healed_items_map.entry(heal_type).or_insert(0) += 1; - self.last_heal_activity = SystemTime::now() - .duration_since(UNIX_EPOCH) - .expect("Time went backwards") - .as_secs(); + pub async fn count_healed(&self, heal_type: HealItemType) { + *self.healed_items_map.write().await.entry(heal_type).or_insert(0) += 1; + *self.last_heal_activity.write().await = SystemTime::now(); } async fn is_quitting(&self) -> bool { - let mut w = self.rx.write().await; - if w.try_recv().is_ok() { + if let Ok(true) = self.rx.has_changed() { + info!("quited"); return true; } false @@ -272,8 +269,7 @@ impl HealSequence { } async fn stop(&self) { - let w = self.tx.write().await; - let _ = w.send(true); + let _ = self.tx.send(true); } async fn push_heal_result_item(&self, r: &HealResultItem) -> Result<()> { @@ -316,19 +312,20 @@ impl HealSequence { Ok(()) } - pub async fn queue_heal_task(&mut self, source: HealSource, heal_type: HealItemType) -> Result<()> { + pub async fn queue_heal_task(&self, source: HealSource, heal_type: HealItemType) -> Result<()> { let mut task = HealTask::new(&source.bucket, &source.object, &source.version_id, &self.setting); + info!("queue_heal_task, {:?}", task); if let Some(opts) = source.opts { task.opts = opts; } else { task.opts.scan_mode = HEAL_UNKNOWN_SCAN; } - self.count_scanned(heal_type.clone()); + self.count_scanned(heal_type.clone()).await; if source.no_wait { let task_str = format!("{:?}", task); - if GLOBAL_BackgroundHealRoutine.read().await.tasks_tx.try_send(task).is_ok() { + if GLOBAL_BackgroundHealRoutine.tasks_tx.try_send(task).is_ok() { info!("Task in the queue: {:?}", task_str); } return Ok(()); @@ -338,8 +335,10 @@ impl HealSequence { task.resp_tx = Some(resp_tx); let task_str = format!("{:?}", task); - if GLOBAL_BackgroundHealRoutine.read().await.tasks_tx.try_send(task).is_ok() { + if GLOBAL_BackgroundHealRoutine.tasks_tx.try_send(task).is_ok() { info!("Task in the queue: {:?}", task_str); + } else { + error!("push task to queue failed"); } let count_ok_drives = |drivers: &[HealDriveInfo]| { let mut count = 0; @@ -354,9 +353,9 @@ impl HealSequence { match resp_rx.recv().await { Some(mut res) => { if res.err.is_none() { - self.count_healed(heal_type.clone()); + self.count_healed(heal_type.clone()).await; } else { - self.count_failed(heal_type.clone()); + self.count_failed(heal_type.clone()).await; } if !self.report_progress { if let Some(err) = res.err { @@ -382,55 +381,67 @@ impl HealSequence { ); } } + + info!("queue_heal_task, HealResult: {:?}", res); self.push_heal_result_item(&res.result).await } None => Ok(()), } } - async fn heal_disk_meta(h: Arc>) -> Result<()> { + async fn heal_disk_meta(h: Arc) -> Result<()> { HealSequence::heal_rustfs_sys_meta(h, CONFIG_PREFIX).await } - async fn heal_items(h: Arc>, buckets_only: bool) -> Result<()> { - if h.read().await.client_token == *BG_HEALING_UUID { + async fn heal_items(h: Arc, buckets_only: bool) -> Result<()> { + if h.client_token == *BG_HEALING_UUID { return Ok(()); } - Self::heal_disk_meta(h.clone()).await?; - let bucket = h.read().await.bucket.clone(); - Self::heal_bucket(h.clone(), &bucket, buckets_only).await + let bucket = h.bucket.clone(); + let task1 = Self::heal_disk_meta(h.clone()); + let task2 = Self::heal_bucket(h.clone(), &bucket, buckets_only); + let results = join!(task1, task2); + results.0?; + results.1?; + + Ok(()) } - async fn traverse_and_heal(h: Arc>) { + async fn traverse_and_heal(h: Arc) { let buckets_only = false; let result = match Self::heal_items(h.clone(), buckets_only).await { Ok(_) => None, Err(err) => Some(err), }; - let _ = h.read().await.traverse_and_heal_done_tx.read().await.send(result).await; + let _ = h.traverse_and_heal_done_tx.read().await.send(result).await; } - async fn heal_rustfs_sys_meta(h: Arc>, meta_prefix: &str) -> Result<()> { - let Some(store) = new_object_layer_fn() else { return Err(Error::msg("errServerNotInitialized")) }; - let setting = h.read().await.setting; + async fn heal_rustfs_sys_meta(h: Arc, meta_prefix: &str) -> Result<()> { + info!("heal_rustfs_sys_meta, h: {:?}", h); + 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(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()))), + }; + let setting = h.setting; store .heal_objects(RUSTFS_META_BUCKET, meta_prefix, &setting, h.clone(), true) .await } async fn is_done(&self) -> bool { - let mut rx_w = self.rx.write().await; - if let Ok(true) = rx_w.recv().await { + if let Ok(true) = self.rx.has_changed() { return true; } false } - pub async fn heal_bucket(hs: Arc>, bucket: &str, bucket_only: bool) -> Result<()> { + pub async fn heal_bucket(hs: Arc, bucket: &str, bucket_only: bool) -> Result<()> { + info!("heal_bucket, hs: {:?}", hs); let (object, setting) = { - let mut hs_w = hs.write().await; - hs_w.queue_heal_task( + hs.queue_heal_task( HealSource { bucket: bucket.to_string(), ..Default::default() @@ -443,37 +454,38 @@ impl HealSequence { return Ok(()); } - if !hs_w.setting.recursive { - if !hs_w.object.is_empty() { - HealSequence::heal_object(hs.clone(), bucket, &hs_w.object, "", hs_w.setting.scan_mode).await?; + if !hs.setting.recursive { + if !hs.object.is_empty() { + HealSequence::heal_object(hs.clone(), bucket, &hs.object, "", hs.setting.scan_mode).await?; } return Ok(()); } - (hs_w.object.clone(), hs_w.setting) + (hs.object.clone(), hs.setting) }; let Some(store) = new_object_layer_fn() else { return Err(Error::msg("errServerNotInitialized")) }; store.heal_objects(bucket, &object, &setting, hs.clone(), false).await } pub async fn heal_object( - hs: Arc>, + hs: Arc, bucket: &str, object: &str, version_id: &str, _scan_mode: HealScanMode, ) -> Result<()> { - let mut hs_w = hs.write().await; - if hs_w.is_quitting().await { + info!("heal_object"); + if hs.is_quitting().await { + info!("heal_object hs is quitting"); return Err(Error::from_string(ERR_HEAL_STOP_SIGNALLED)); } - let setting = hs_w.setting; - hs_w.queue_heal_task( + info!("will queue task"); + hs.queue_heal_task( HealSource { bucket: bucket.to_string(), object: object.to_string(), version_id: version_id.to_string(), - opts: Some(setting), + opts: Some(hs.setting), ..Default::default() }, HEAL_ITEM_OBJECT.to_string(), @@ -484,18 +496,17 @@ impl HealSequence { } pub async fn heal_meta_object( - hs: Arc>, + hs: Arc, bucket: &str, object: &str, version_id: &str, _scan_mode: HealScanMode, ) -> Result<()> { - let mut hs_w = hs.write().await; - if hs_w.is_quitting().await { + if hs.is_quitting().await { return Err(Error::from_string(ERR_HEAL_STOP_SIGNALLED)); } - hs_w.queue_heal_task( + hs.queue_heal_task( HealSource { bucket: bucket.to_string(), object: object.to_string(), @@ -510,10 +521,9 @@ impl HealSequence { } } -pub async fn heal_sequence_start(h: Arc>) { - let r = h.read().await; +pub async fn heal_sequence_start(h: Arc) { { - let mut current_status_w = r.current_status.write().await; + let mut current_status_w = h.current_status.write().await; current_status_w.summary = HEAL_RUNNING_STATUS.to_string(); current_status_w.start_time = SystemTime::now() .duration_since(UNIX_EPOCH) @@ -527,16 +537,15 @@ pub async fn heal_sequence_start(h: Arc>) { }); let h_clone_1 = h.clone(); - let mut x = r.traverse_and_heal_done_rx.write().await; + let mut x = h.traverse_and_heal_done_rx.write().await; select! { - _ = r.is_done() => { - *(r.end_time.write().await) = SystemTime::now(); - let mut current_status_w = r.current_status.write().await; + _ = h.is_done() => { + *(h.end_time.write().await) = SystemTime::now(); + let mut current_status_w = h.current_status.write().await; current_status_w.summary = HEAL_FINISHED_STATUS.to_string(); spawn(async move { - let binding = h_clone_1.read().await; - let mut rx_w = binding.traverse_and_heal_done_rx.write().await; + let mut rx_w = h_clone_1.traverse_and_heal_done_rx.write().await; rx_w.recv().await; }); } @@ -544,12 +553,12 @@ pub async fn heal_sequence_start(h: Arc>) { if let Some(err) = result { match err { Some(err) => { - let mut current_status_w = r.current_status.write().await; + let mut current_status_w = h.current_status.write().await; (current_status_w).summary = HEAL_STOPPED_STATUS.to_string(); (current_status_w).failure_detail = err.to_string(); }, None => { - let mut current_status_w = r.current_status.write().await; + let mut current_status_w = h.current_status.write().await; (current_status_w).summary = HEAL_FINISHED_STATUS.to_string(); } } @@ -563,7 +572,7 @@ pub async fn heal_sequence_start(h: Arc>) { pub struct AllHealState { mu: RwLock, - heal_seq_map: HashMap>>, + heal_seq_map: HashMap>, heal_local_disks: HashMap, heal_status: HashMap, } @@ -664,8 +673,7 @@ impl AllHealState { let mut keys_to_reomve = Vec::new(); for (k, v) in self.heal_seq_map.iter() { - let r = v.read().await; - if r.has_ended().await && now.duration_since(*(r.end_time.read().await)).unwrap() > KEEP_HEAL_SEQ_STATE_DURATION { + if v.has_ended().await && now.duration_since(*(v.end_time.read().await)).unwrap() > KEEP_HEAL_SEQ_STATE_DURATION { keys_to_reomve.push(k.clone()) } } @@ -674,12 +682,11 @@ impl AllHealState { } } - pub async fn get_heal_sequence_by_token(&self, token: &str) -> (Option>>, bool) { + pub async fn get_heal_sequence_by_token(&self, token: &str) -> (Option>, bool) { let _ = self.mu.read().await; for v in self.heal_seq_map.values() { - let r = v.read().await; - if r.client_token == token { + if v.client_token == token { return (Some(v.clone()), true); } } @@ -687,7 +694,7 @@ impl AllHealState { (None, false) } - pub async fn get_heal_sequence(&self, path: &str) -> Option>> { + pub async fn get_heal_sequence(&self, path: &str) -> Option> { let _ = self.mu.read().await; self.heal_seq_map.get(path).cloned() @@ -696,7 +703,6 @@ impl AllHealState { pub async fn stop_heal_sequence(&mut self, path: &str) -> Result> { let mut hsp = HealStopSuccess::default(); if let Some(he) = self.get_heal_sequence(path).await { - let he = he.read().await; let client_token = he.client_token.clone(); if *GLOBAL_IsDistErasure.read().await { // TODO: proxy @@ -736,22 +742,21 @@ impl AllHealState { // `keepHealSeqStateDuration`. This function also launches a // background routine to clean up heal results after the // aforementioned duration. - pub async fn launch_new_heal_sequence(&mut self, heal_sequence: Arc>) -> Result> { - let r = heal_sequence.read().await; - let path = Path::new(&r.bucket).join(r.object.clone()); + pub async fn launch_new_heal_sequence(&mut self, heal_sequence: Arc) -> Result> { + let path = Path::new(&heal_sequence.bucket).join(heal_sequence.object.clone()); let path_s = path.to_str().unwrap(); - if r.force_started { + if heal_sequence.force_started { self.stop_heal_sequence(path_s).await?; } else if let Some(hs) = self.get_heal_sequence(path_s).await { - if !hs.read().await.has_ended().await { - return Err(Error::from_string(format!("Heal is already running on the given path (use force-start option to stop and start afresh). The heal was started by IP {} at {:?}, token is {}", r.client_address, r.start_time, r.client_token))); + if !hs.has_ended().await { + return Err(Error::from_string(format!("Heal is already running on the given path (use force-start option to stop and start afresh). The heal was started by IP {} at {:?}, token is {}", heal_sequence.client_address, heal_sequence.start_time, heal_sequence.client_token))); } } let _ = self.mu.write().await; for (k, v) in self.heal_seq_map.iter() { - if !v.read().await.has_ended().await && (has_profix(k, path_s) || has_profix(path_s, k)) { + if !v.has_ended().await && (has_profix(k, path_s) || has_profix(path_s, k)) { return Err(Error::from_string(format!( "The provided heal sequence path overlaps with an existing heal path: {}", k @@ -761,12 +766,12 @@ impl AllHealState { self.heal_seq_map.insert(path_s.to_string(), heal_sequence.clone()); - let client_token = r.client_token.clone(); + let client_token = heal_sequence.client_token.clone(); if *GLOBAL_IsDistErasure.read().await { // TODO: proxy } - if r.client_token == BG_HEALING_UUID { + if heal_sequence.client_token == BG_HEALING_UUID { // For background heal do nothing, do not spawn an unnecessary goroutine. } else { let heal_sequence_clone = heal_sequence.clone(); @@ -777,8 +782,8 @@ impl AllHealState { let b = serde_json::to_vec(&HealStartSuccess { client_token, - client_address: r.client_address.clone(), - start_time: r.start_time, + client_address: heal_sequence.client_address.clone(), + start_time: heal_sequence.start_time, })?; Ok(b) } diff --git a/ecstore/src/peer.rs b/ecstore/src/peer.rs index 369b47dfb..570ed0939 100644 --- a/ecstore/src/peer.rs +++ b/ecstore/src/peer.rs @@ -8,7 +8,7 @@ use regex::Regex; use std::{collections::HashMap, fmt::Debug, sync::Arc}; use tokio::sync::RwLock; use tonic::Request; -use tracing::warn; +use tracing::{error, info, warn}; use crate::disk::error::{is_all_buckets_not_found, is_all_not_found}; use crate::disk::{DiskAPI, DiskStore}; @@ -102,15 +102,17 @@ impl S3PeerSys { pool_errs.push(reduce_write_quorum_errs(&per_pool_errs, &bucket_op_ignored_errs(), qu)); } + error!("found pool errs: {:?}", pool_errs); + if !opts.recreate { - opts.remove = is_all_not_found(&pool_errs); - opts.recursive = !opts.remove; + opts.remove = is_all_buckets_not_found(&pool_errs); + opts.recreate = !opts.remove; } let mut futures = Vec::new(); let heal_bucket_results = Arc::new(RwLock::new(vec![HealResultItem::default(); self.clients.len()])); for (idx, client) in self.clients.iter().enumerate() { - let opts_clone = opts; + let opts_clone = opts.clone(); let heal_bucket_results_clone = heal_bucket_results.clone(); futures.push(async move { match client.heal_bucket(bucket, &opts_clone).await { @@ -723,14 +725,16 @@ pub async fn heal_bucket_local(bucket: &str, opts: &HealOpts) -> Result { + info!("will call delete_volume, volume: {}", bucket); let _ = disk.delete_volume(&bucket).await; None } @@ -752,6 +756,7 @@ pub async fn heal_bucket_local(bucket: &str, opts: &HealOpts) -> Result { as_clone.write().await[idx] = DRIVE_STATE_OK.to_string(); diff --git a/ecstore/src/sets.rs b/ecstore/src/sets.rs index 8dd190839..b68e3283d 100644 --- a/ecstore/src/sets.rs +++ b/ecstore/src/sets.rs @@ -662,7 +662,7 @@ impl StorageAPI for Sets { _bucket: &str, _prefix: &str, _opts: &HealOpts, - _hs: Arc>, + _hs: Arc, _is_meta: bool, ) -> Result<()> { unimplemented!() diff --git a/ecstore/src/store.rs b/ecstore/src/store.rs index 2ae73529d..18e19f30e 100644 --- a/ecstore/src/store.rs +++ b/ecstore/src/store.rs @@ -1953,6 +1953,7 @@ impl StorageAPI for ECStore { version_id: &str, opts: &HealOpts, ) -> Result<(HealResultItem, Option)> { + info!("ECStore heal_object"); let object = utils::path::encode_dir_object(object); let errs = Arc::new(RwLock::new(vec![None; self.pools.len()])); let results = Arc::new(RwLock::new(vec![HealResultItem::default(); self.pools.len()])); @@ -2009,9 +2010,10 @@ impl StorageAPI for ECStore { bucket: &str, prefix: &str, opts: &HealOpts, - hs: Arc>, + hs: Arc, is_meta: bool, ) -> Result<()> { + info!("heal objects"); let opts_clone = *opts; let heal_entry: HealEntryFn = Arc::new(move |bucket: String, entry: MetaCacheEntry, scan_mode: HealScanMode| { let opts_clone = opts_clone; diff --git a/ecstore/src/store_api.rs b/ecstore/src/store_api.rs index 75cb4133b..bf8505726 100644 --- a/ecstore/src/store_api.rs +++ b/ecstore/src/store_api.rs @@ -14,7 +14,6 @@ use serde::{Deserialize, Serialize}; use std::collections::HashMap; use std::sync::Arc; use time::OffsetDateTime; -use tokio::sync::RwLock; use uuid::Uuid; pub const ERASURE_ALGORITHM: &str = "rs-vandermonde"; @@ -292,7 +291,7 @@ pub struct ObjectPartInfo { // } // } -#[derive(Serialize, Deserialize)] +#[derive(Default, Serialize, Deserialize)] pub struct RawFileInfo { pub buf: Vec, } @@ -1028,7 +1027,7 @@ pub trait StorageAPI: ObjectIO { bucket: &str, prefix: &str, opts: &HealOpts, - hs: Arc>, + hs: Arc, is_meta: bool, ) -> Result<()>; async fn get_pool_and_set(&self, id: &str) -> Result<(Option, Option, Option)>; diff --git a/rustfs/src/main.rs b/rustfs/src/main.rs index edbfa7533..92460bd24 100644 --- a/rustfs/src/main.rs +++ b/rustfs/src/main.rs @@ -192,9 +192,9 @@ async fn run(opt: config::Opt) -> Result<()> { })?; warn!(" init store success!"); // init scanner - // init_data_scanner().await; - // // init auto heal - // init_auto_heal().await; + init_data_scanner().await; + // init auto heal + init_auto_heal().await; info!("server was started");