heal admin api(2)

Signed-off-by: mujunxiang <1948535941@qq.com>
This commit is contained in:
mujunxiang
2024-11-23 14:18:35 +08:00
parent a421b5be2b
commit 6ddfc80943
12 changed files with 167 additions and 150 deletions
+9 -6
View File
@@ -419,15 +419,18 @@ pub fn convert_access_error(e: io::Error, per_err: DiskError) -> Error {
}
pub fn is_all_not_found(errs: &[Option<Error>]) -> bool {
for err in errs.iter().flatten() {
if let Some(err) = err.downcast_ref::<DiskError>() {
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::<DiskError>() {
match err {
DiskError::FileNotFound | DiskError::VolumeNotFound | &DiskError::FileVersionNotFound => {
continue;
}
_ => return false,
}
_ => return false,
}
}
return false;
}
!errs.is_empty()
+2
View File
@@ -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
+1
View File
@@ -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,
+1 -1
View File
@@ -28,7 +28,7 @@ lazy_static! {
pub static ref GLOBAL_LOCAL_DISK_SET_DRIVES: Arc<RwLock<TypeLocalDiskSetDrives>> = Arc::new(RwLock::new(Vec::new()));
pub static ref GLOBAL_Endpoints: OnceLock<EndpointServerPools> = OnceLock::new();
pub static ref GLOBAL_RootDiskThreshold: RwLock<u64> = RwLock::new(0);
pub static ref GLOBAL_BackgroundHealRoutine: Arc<RwLock<HealRoutine>> = HealRoutine::new();
pub static ref GLOBAL_BackgroundHealRoutine: Arc<HealRoutine> = HealRoutine::new();
pub static ref GLOBAL_BackgroundHealState: Arc<RwLock<AllHealState>> = AllHealState::new(false);
static ref globalDeploymentIDPtr: RwLock<Uuid> = RwLock::new(Uuid::nil());
}
+19 -13
View File
@@ -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<Error>,
@@ -325,12 +326,12 @@ pub struct HealResult {
pub struct HealRoutine {
pub tasks_tx: Sender<HealTask>,
tasks_rx: Receiver<HealTask>,
tasks_rx: RwLock<Receiver<HealTask>>,
workers: usize,
}
impl HealRoutine {
pub fn new() -> Arc<RwLock<Self>> {
pub fn new() -> Arc<Self> {
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::<usize>() {
@@ -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<RwLock<HealSequence>>) {
pub async fn add_worker(&self, bgseq: Arc<HealSequence>) {
loop {
let mut d_res = HealResultItem::default();
let d_err: Option<Error>;
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;
},
}
}
}
-6
View File
@@ -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(),
+116 -111
View File
@@ -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<HealOpts>,
}
#[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<RwLock<HealSequenceStatus>>,
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<ItemsMap>,
pub healed_items_map: RwLock<ItemsMap>,
pub heal_failed_items_map: RwLock<ItemsMap>,
pub last_heal_activity: RwLock<SystemTime>,
traverse_and_heal_done_tx: Arc<RwLock<M_Sender<Option<Error>>>>,
traverse_and_heal_done_rx: Arc<RwLock<M_Receiver<Option<Error>>>>,
tx: Arc<RwLock<Sender<bool>>>,
rx: Arc<RwLock<Receiver<bool>>>,
tx: W_Sender<bool>,
rx: W_Receiver<bool>,
}
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<RwLock<HealSequence>>) -> Result<()> {
async fn heal_disk_meta(h: Arc<HealSequence>) -> Result<()> {
HealSequence::heal_rustfs_sys_meta(h, CONFIG_PREFIX).await
}
async fn heal_items(h: Arc<RwLock<HealSequence>>, buckets_only: bool) -> Result<()> {
if h.read().await.client_token == *BG_HEALING_UUID {
async fn heal_items(h: Arc<HealSequence>, 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<RwLock<HealSequence>>) {
async fn traverse_and_heal(h: Arc<HealSequence>) {
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<RwLock<HealSequence>>, 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<HealSequence>, 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<RwLock<HealSequence>>, bucket: &str, bucket_only: bool) -> Result<()> {
pub async fn heal_bucket(hs: Arc<HealSequence>, 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<RwLock<HealSequence>>,
hs: Arc<HealSequence>,
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<RwLock<HealSequence>>,
hs: Arc<HealSequence>,
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<RwLock<HealSequence>>) {
let r = h.read().await;
pub async fn heal_sequence_start(h: Arc<HealSequence>) {
{
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<RwLock<HealSequence>>) {
});
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<RwLock<HealSequence>>) {
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<RwLock<HealSequence>>) {
pub struct AllHealState {
mu: RwLock<bool>,
heal_seq_map: HashMap<String, Arc<RwLock<HealSequence>>>,
heal_seq_map: HashMap<String, Arc<HealSequence>>,
heal_local_disks: HashMap<Endpoint, bool>,
heal_status: HashMap<String, HealingTracker>,
}
@@ -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<Arc<RwLock<HealSequence>>>, bool) {
pub async fn get_heal_sequence_by_token(&self, token: &str) -> (Option<Arc<HealSequence>>, 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<Arc<RwLock<HealSequence>>> {
pub async fn get_heal_sequence(&self, path: &str) -> Option<Arc<HealSequence>> {
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<Vec<u8>> {
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<RwLock<HealSequence>>) -> Result<Vec<u8>> {
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<HealSequence>) -> Result<Vec<u8>> {
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)
}
+10 -5
View File
@@ -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<HealResu
});
}
if !is_rustfs_meta_bucket_name(bucket) && !is_all_buckets_not_found(&errs) && opts.remove {
if opts.remove && !is_rustfs_meta_bucket_name(bucket) && !is_all_buckets_not_found(&errs) {
let mut futures = Vec::new();
for disk in disks.iter() {
let disk = disk.clone();
let bucket = bucket.to_string();
info!("heal_bucket_local, errs: {:?}, opts: {:?}", errs, opts);
futures.push(async move {
match disk {
Some(disk) => {
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<HealResu
let errs_clone = errs.clone();
futures.push(async move {
if bs_clone.read().await[idx] == DRIVE_STATE_MISSING {
info!("bucket not find, will recreate");
match disk.as_ref().unwrap().make_volume(&bucket).await {
Ok(_) => {
as_clone.write().await[idx] = DRIVE_STATE_OK.to_string();
+1 -1
View File
@@ -662,7 +662,7 @@ impl StorageAPI for Sets {
_bucket: &str,
_prefix: &str,
_opts: &HealOpts,
_hs: Arc<RwLock<HealSequence>>,
_hs: Arc<HealSequence>,
_is_meta: bool,
) -> Result<()> {
unimplemented!()
+3 -1
View File
@@ -1953,6 +1953,7 @@ impl StorageAPI for ECStore {
version_id: &str,
opts: &HealOpts,
) -> Result<(HealResultItem, Option<Error>)> {
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<RwLock<HealSequence>>,
hs: Arc<HealSequence>,
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;
+2 -3
View File
@@ -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<u8>,
}
@@ -1028,7 +1027,7 @@ pub trait StorageAPI: ObjectIO {
bucket: &str,
prefix: &str,
opts: &HealOpts,
hs: Arc<RwLock<HealSequence>>,
hs: Arc<HealSequence>,
is_meta: bool,
) -> Result<()>;
async fn get_pool_and_set(&self, id: &str) -> Result<(Option<usize>, Option<usize>, Option<usize>)>;
+3 -3
View File
@@ -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");