diff --git a/crates/heal/src/heal/channel.rs b/crates/heal/src/heal/channel.rs index f0443d19a..a06c1754e 100644 --- a/crates/heal/src/heal/channel.rs +++ b/crates/heal/src/heal/channel.rs @@ -267,6 +267,14 @@ impl HealChannelProcessor { status: HealTaskStatus::Pending | HealTaskStatus::Running, result_items, }) => ("running".to_string(), None, result_items), + Ok(HealTaskReport { + status: HealTaskStatus::Retrying { error, retry_attempt }, + result_items, + }) => ( + "running".to_string(), + Some(format!("heal task retrying after recoverable failure, attempt {retry_attempt}: {error}")), + result_items, + ), Ok(HealTaskReport { status: HealTaskStatus::Completed, result_items, diff --git a/crates/heal/src/heal/manager.rs b/crates/heal/src/heal/manager.rs index da35861ce..985965fd8 100644 --- a/crates/heal/src/heal/manager.rs +++ b/crates/heal/src/heal/manager.rs @@ -31,7 +31,7 @@ use std::{ }; use tokio::{ sync::{Mutex, Notify, RwLock}, - time::interval, + time::{interval, sleep}, }; use tokio_util::sync::CancellationToken; use tracing::{debug, error, info, warn}; @@ -47,6 +47,8 @@ const EVENT_HEAL_MANAGER_STATE: &str = "heal_manager_state"; const EVENT_HEAL_QUEUE_ADMISSION: &str = "heal_queue_admission"; const EVENT_HEAL_SCHEDULER_STATE: &str = "heal_scheduler_state"; const EVENT_HEAL_QUEUE_STATE: &str = "heal_queue_state"; +const MAX_RECOVERABLE_HEAL_RETRIES: u32 = 3; +const MAX_RECOVERABLE_HEAL_RETRY_DELAY: Duration = Duration::from_secs(30); /// Priority queue wrapper for heal requests /// Uses BinaryHeap for priority-based ordering while maintaining FIFO for same-priority items @@ -109,6 +111,18 @@ struct CompletedHealStatus { completed_at: SystemTime, } +#[derive(Debug, Clone)] +struct HealTaskAlias { + task_id: String, +} + +#[derive(Debug, Clone)] +struct RetryingHeal { + request: HealRequest, + error: String, + cancel_token: CancellationToken, +} + #[derive(Debug, Clone)] pub struct HealTaskReport { pub status: HealTaskStatus, @@ -330,7 +344,11 @@ impl PriorityHealQueue { /// Create a deduplication key from a heal request fn make_dedup_key(request: &HealRequest) -> String { - match &request.heal_type { + Self::make_dedup_key_for_type(&request.heal_type) + } + + fn make_dedup_key_for_type(heal_type: &HealType) -> String { + match heal_type { HealType::Object { bucket, object, @@ -393,6 +411,12 @@ impl PriorityHealQueue { .any(|item| item.request.id == request_id && heal_type_matches_path(&item.request.heal_type, heal_path)) } + fn request_for_dedup_key(&self, key: &str) -> Option<&HealRequest> { + self.heap + .iter() + .find_map(|item| (Self::make_dedup_key(&item.request) == key).then_some(&item.request)) + } + fn contains_matching(&self, mut matches: F) -> bool where F: FnMut(&HealRequest) -> bool, @@ -418,25 +442,34 @@ impl PriorityHealQueue { removed } - fn remove_matching(&mut self, mut should_remove: F) -> usize + fn remove_matching(&mut self, mut should_remove: F) -> Vec where F: FnMut(&HealRequest) -> bool, { let mut retained = BinaryHeap::new(); - let mut removed_count = 0; + let mut removed = Vec::new(); while let Some(item) = self.heap.pop() { if should_remove(&item.request) { let key = Self::make_dedup_key(&item.request); Self::decrement_or_remove_dedup_key(&mut self.dedup_keys, &key); - removed_count += 1; + removed.push(item.request); } else { retained.push(item); } } self.heap = retained; - removed_count + removed + } +} + +impl RetryingHeal { + fn status(&self) -> HealTaskStatus { + HealTaskStatus::Retrying { + error: self.error.clone(), + retry_attempt: self.request.retry_attempts, + } } } @@ -464,6 +497,86 @@ fn publish_heal_queue_length(queue: &PriorityHealQueue) { crate::set_heal_queue_length(queue.len()); } +fn active_heals_contains_dedup_key(active_heals: &HashMap>, key: &str) -> bool { + active_heal_for_dedup_key(active_heals, key).is_some() +} + +fn active_heal_for_dedup_key(active_heals: &HashMap>, key: &str) -> Option<(String, HealType)> { + active_heals + .iter() + .find(|(_, task)| PriorityHealQueue::make_dedup_key_for_type(&task.heal_type) == key) + .map(|(task_id, task)| (task_id.clone(), task.heal_type.clone())) +} + +fn retrying_heal_for_dedup_key(retrying_heals: &HashMap, key: &str) -> Option<(String, HealType)> { + retrying_heals + .iter() + .find(|(_, retrying)| PriorityHealQueue::make_dedup_key(&retrying.request) == key) + .map(|(task_id, retrying)| (task_id.clone(), retrying.request.heal_type.clone())) +} + +fn completed_status_is_retrying(status: &HealTaskStatus) -> bool { + matches!(status, HealTaskStatus::Retrying { .. }) +} + +fn retry_request_for_result(task: &HealTask, result: &Result<()>) -> Option<(HealRequest, Duration, String)> { + let Err(err) = result else { + return None; + }; + if task.retry_attempts >= MAX_RECOVERABLE_HEAL_RETRIES { + return None; + } + + let error = err.to_string(); + if !is_recoverable_heal_error(err, &error) { + return None; + } + + let request = task.retry_request(); + let delay = recoverable_heal_retry_delay(request.retry_attempts); + Some((request, delay, error)) +} + +fn recoverable_heal_retry_delay(retry_attempt: u32) -> Duration { + let retry_attempt = retry_attempt.clamp(1, 5); + let delay = Duration::from_secs(2_u64.saturating_pow(retry_attempt)); + delay.min(MAX_RECOVERABLE_HEAL_RETRY_DELAY) +} + +fn is_recoverable_heal_error(err: &Error, error: &str) -> bool { + match err { + Error::TaskCancelled => false, + Error::TaskTimeout | Error::TransientSkip { .. } => true, + Error::TaskExecutionFailed { .. } + | Error::Storage(_) + | Error::Disk(_) + | Error::Io(_) + | Error::IO(_) + | Error::Anyhow(_) + | Error::Other(_) => is_recoverable_heal_error_message(error), + _ => false, + } +} + +fn is_recoverable_heal_error_message(error: &str) -> bool { + let error = error.to_ascii_lowercase(); + [ + "failed to acquire read lock", + "lock acquisition failed", + "lock acquisition timeout", + "remote lock rpc timed out", + "deadline has elapsed", + "timed out", + "transport error", + "network error", + "connection refused", + "operation canceled", + "quorum not reached", + ] + .iter() + .any(|pattern| error.contains(pattern)) +} + /// Heal config #[derive(Debug, Clone)] pub struct HealConfig { @@ -578,6 +691,10 @@ pub struct HealManager { heal_queue: Arc>, /// Recently completed heal statuses retained for status queries. completed_heals: Arc>>, + /// Client tokens merged into an existing task id. + task_aliases: Arc>>, + /// Heal tasks waiting for a retry backoff to expire. + retrying_heals: Arc>>, /// Storage layer interface storage: Arc, /// Cancel token @@ -588,6 +705,18 @@ pub struct HealManager { notify: Arc, } +struct HealQueueContext<'a> { + heal_queue: &'a Arc>, + active_heals: &'a Arc>>>, + completed_heals: &'a Arc>>, + retrying_heals: &'a Arc>>, + config: &'a Arc>, + statistics: &'a Arc>, + storage: &'a Arc, + notify: &'a Arc, + cancel_token: &'a CancellationToken, +} + impl HealManager { fn classify_full_admission(request: &HealRequest, config: &HealConfig) -> HealAdmissionResult { if request.priority == HealPriority::Low && config.low_priority_drop_when_full { @@ -597,6 +726,43 @@ impl HealManager { } } + fn duplicate_admission_for_request(request: &HealRequest, config: &HealConfig) -> HealAdmissionResult { + if request.priority == HealPriority::Low && !config.low_priority_merge_enable { + HealAdmissionResult::Dropped(HealAdmissionDropReason::PolicyDropped) + } else { + HealAdmissionResult::Merged + } + } + + async fn insert_task_alias(&self, alias_id: &str, task_id: &str) { + if alias_id == task_id { + return; + } + + self.task_aliases.lock().await.insert( + alias_id.to_string(), + HealTaskAlias { + task_id: task_id.to_string(), + }, + ); + } + + async fn canonical_task_id(&self, task_id: &str) -> String { + self.task_aliases + .lock() + .await + .get(task_id) + .map(|alias| alias.task_id.clone()) + .unwrap_or_else(|| task_id.to_string()) + } + + async fn remove_aliases_for_task(&self, task_id: &str) { + self.task_aliases + .lock() + .await + .retain(|alias_id, alias| alias_id != task_id && alias.task_id != task_id); + } + /// Create new HealManager pub fn new(storage: Arc, config: Option) -> Self { let config = config.unwrap_or_default(); @@ -606,6 +772,8 @@ impl HealManager { active_heals: Arc::new(Mutex::new(HashMap::new())), heal_queue: Arc::new(Mutex::new(PriorityHealQueue::new())), completed_heals: Arc::new(Mutex::new(HashMap::new())), + task_aliases: Arc::new(Mutex::new(HashMap::new())), + retrying_heals: Arc::new(Mutex::new(HashMap::new())), storage, cancel_token: CancellationToken::new(), statistics: Arc::new(RwLock::new(HealStatistics::new())), @@ -689,6 +857,8 @@ impl HealManager { active_heals.clear(); publish_active_heal_count(&active_heals); self.completed_heals.lock().await.clear(); + self.task_aliases.lock().await.clear(); + self.retrying_heals.lock().await.clear(); crate::set_heal_queue_length(0); // update state @@ -709,27 +879,119 @@ impl HealManager { /// Submit heal request pub async fn submit_heal_request(&self, request: HealRequest) -> Result { let config = self.config.read().await; + let dedup_key = PriorityHealQueue::make_dedup_key(&request); + + // Keep this lock order aligned with the scheduler so a request cannot + // slip between queued and active states while duplicate admission runs. + let active_duplicate = { + let active_heals = self.active_heals.lock().await; + (!request.force_start) + .then(|| active_heal_for_dedup_key(&active_heals, &dedup_key)) + .flatten() + }; + if let Some((merged_task_id, _)) = active_duplicate { + let admission = Self::duplicate_admission_for_request(&request, &config); + + match admission { + HealAdmissionResult::Merged => { + self.insert_task_alias(&request.id, &merged_task_id).await; + info!( + target: "rustfs::heal::manager", + event = EVENT_HEAL_QUEUE_ADMISSION, + component = LOG_COMPONENT_HEAL, + subsystem = LOG_SUBSYSTEM_MANAGER, + request_id = %request.id, + merged_task_id = %merged_task_id, + priority = ?request.priority, + result = "merged_active_duplicate", + "Heal queue admission decided" + ); + } + HealAdmissionResult::Dropped(reason) => { + warn!( + target: "rustfs::heal::manager", + event = EVENT_HEAL_QUEUE_ADMISSION, + component = LOG_COMPONENT_HEAL, + subsystem = LOG_SUBSYSTEM_MANAGER, + request_id = %request.id, + priority = ?request.priority, + reason = reason.as_str(), + result = "dropped_active_duplicate", + "Heal queue admission decided" + ); + } + HealAdmissionResult::Accepted | HealAdmissionResult::Full => {} + } + + return Ok(admission); + } + + let retrying_duplicate = { + let retrying_heals = self.retrying_heals.lock().await; + (!request.force_start) + .then(|| retrying_heal_for_dedup_key(&retrying_heals, &dedup_key)) + .flatten() + }; + if let Some((merged_task_id, _)) = retrying_duplicate { + let admission = Self::duplicate_admission_for_request(&request, &config); + + match admission { + HealAdmissionResult::Merged => { + self.insert_task_alias(&request.id, &merged_task_id).await; + info!( + target: "rustfs::heal::manager", + event = EVENT_HEAL_QUEUE_ADMISSION, + component = LOG_COMPONENT_HEAL, + subsystem = LOG_SUBSYSTEM_MANAGER, + request_id = %request.id, + merged_task_id = %merged_task_id, + priority = ?request.priority, + result = "merged_retrying_duplicate", + "Heal queue admission decided" + ); + } + HealAdmissionResult::Dropped(reason) => { + warn!( + target: "rustfs::heal::manager", + event = EVENT_HEAL_QUEUE_ADMISSION, + component = LOG_COMPONENT_HEAL, + subsystem = LOG_SUBSYSTEM_MANAGER, + request_id = %request.id, + priority = ?request.priority, + reason = reason.as_str(), + result = "dropped_retrying_duplicate", + "Heal queue admission decided" + ); + } + HealAdmissionResult::Accepted | HealAdmissionResult::Full => {} + } + + return Ok(admission); + } + let mut queue = self.heal_queue.lock().await; let queue_len = queue.len(); publish_heal_queue_length(&queue); let queue_capacity = config.queue_size; - if !request.force_start && queue.contains_key(&request) { - let admission = if request.priority == HealPriority::Low && !config.low_priority_merge_enable { - HealAdmissionResult::Dropped(HealAdmissionDropReason::PolicyDropped) - } else { - HealAdmissionResult::Merged - }; + let queued_duplicate = (!request.force_start) + .then(|| queue.request_for_dedup_key(&dedup_key).map(|queued| queued.id.clone())) + .flatten(); + if let Some(merged_task_id) = queued_duplicate { + let admission = Self::duplicate_admission_for_request(&request, &config); + drop(queue); match admission { HealAdmissionResult::Merged => { + self.insert_task_alias(&request.id, &merged_task_id).await; debug!( target: "rustfs::heal::manager", event = EVENT_HEAL_QUEUE_ADMISSION, component = LOG_COMPONENT_HEAL, subsystem = LOG_SUBSYSTEM_MANAGER, request_id = %request.id, + merged_task_id = %merged_task_id, priority = ?request.priority, result = "merged_duplicate", "Heal queue request merged" @@ -897,22 +1159,40 @@ impl HealManager { /// Get task status pub async fn get_task_status(&self, task_id: &str) -> Result { + let canonical_task_id = self.canonical_task_id(task_id).await; { let active_heals = self.active_heals.lock().await; - if let Some(task) = active_heals.get(task_id) { + if let Some(task) = active_heals.get(&canonical_task_id) { return Ok(task.get_status().await); } } + { + let retrying_heals = self.retrying_heals.lock().await; + if let Some(retrying) = retrying_heals.get(&canonical_task_id) { + return Ok(retrying.status()); + } + } + + { + let mut completed_heals = self.completed_heals.lock().await; + prune_completed_heal_statuses(&mut completed_heals); + if let Some(completed) = completed_heals.get(&canonical_task_id) + && completed_status_is_retrying(&completed.status) + { + return Ok(completed.status.clone()); + } + } + let queue = self.heal_queue.lock().await; - if queue.contains_request_id(task_id) { + if queue.contains_request_id(&canonical_task_id) { return Ok(HealTaskStatus::Pending); } drop(queue); let mut completed_heals = self.completed_heals.lock().await; prune_completed_heal_statuses(&mut completed_heals); - if let Some(completed) = completed_heals.get(task_id) { + if let Some(completed) = completed_heals.get(&canonical_task_id) { return Ok(completed.status.clone()); } @@ -922,9 +1202,10 @@ impl HealManager { } pub async fn get_task_report_for_path(&self, heal_path: &str, task_id: &str) -> Result { + let canonical_task_id = self.canonical_task_id(task_id).await; { let active_heals = self.active_heals.lock().await; - if let Some(task) = active_heals.get(task_id) + if let Some(task) = active_heals.get(&canonical_task_id) && heal_type_matches_path(&task.heal_type, heal_path) { return Ok(HealTaskReport { @@ -934,9 +1215,35 @@ impl HealManager { } } + { + let retrying_heals = self.retrying_heals.lock().await; + if let Some(retrying) = retrying_heals.get(&canonical_task_id) + && heal_type_matches_path(&retrying.request.heal_type, heal_path) + { + return Ok(HealTaskReport { + status: retrying.status(), + result_items: Vec::new(), + }); + } + } + + { + let mut completed_heals = self.completed_heals.lock().await; + prune_completed_heal_statuses(&mut completed_heals); + if let Some(completed) = completed_heals.get(&canonical_task_id) + && heal_type_matches_path(&completed.heal_type, heal_path) + && completed_status_is_retrying(&completed.status) + { + return Ok(HealTaskReport { + status: completed.status.clone(), + result_items: completed.result_items.clone(), + }); + } + } + { let queue = self.heal_queue.lock().await; - if queue.contains_request_id_matching_path(task_id, heal_path) { + if queue.contains_request_id_matching_path(&canonical_task_id, heal_path) { return Ok(HealTaskReport { status: HealTaskStatus::Pending, result_items: Vec::new(), @@ -947,7 +1254,7 @@ impl HealManager { { let mut completed_heals = self.completed_heals.lock().await; prune_completed_heal_statuses(&mut completed_heals); - if let Some(completed) = completed_heals.get(task_id) + if let Some(completed) = completed_heals.get(&canonical_task_id) && heal_type_matches_path(&completed.heal_type, heal_path) { return Ok(HealTaskReport { @@ -972,18 +1279,39 @@ impl HealManager { /// treat it as an already-finished sequence. If the path still has a live or /// recently completed task, a different token is invalid for that path. pub async fn get_task_status_for_path(&self, heal_path: &str, task_id: &str) -> Result { + let canonical_task_id = self.canonical_task_id(task_id).await; { let active_heals = self.active_heals.lock().await; - if let Some(task) = active_heals.get(task_id) + if let Some(task) = active_heals.get(&canonical_task_id) && heal_type_matches_path(&task.heal_type, heal_path) { return Ok(task.get_status().await); } } + { + let retrying_heals = self.retrying_heals.lock().await; + if let Some(retrying) = retrying_heals.get(&canonical_task_id) + && heal_type_matches_path(&retrying.request.heal_type, heal_path) + { + return Ok(retrying.status()); + } + } + + { + let mut completed_heals = self.completed_heals.lock().await; + prune_completed_heal_statuses(&mut completed_heals); + if let Some(completed) = completed_heals.get(&canonical_task_id) + && heal_type_matches_path(&completed.heal_type, heal_path) + && completed_status_is_retrying(&completed.status) + { + return Ok(completed.status.clone()); + } + } + { let queue = self.heal_queue.lock().await; - if queue.contains_request_id_matching_path(task_id, heal_path) { + if queue.contains_request_id_matching_path(&canonical_task_id, heal_path) { return Ok(HealTaskStatus::Pending); } } @@ -991,7 +1319,7 @@ impl HealManager { { let mut completed_heals = self.completed_heals.lock().await; prune_completed_heal_statuses(&mut completed_heals); - if let Some(completed) = completed_heals.get(task_id) + if let Some(completed) = completed_heals.get(&canonical_task_id) && heal_type_matches_path(&completed.heal_type, heal_path) { return Ok(completed.status.clone()); @@ -1025,6 +1353,16 @@ impl HealManager { } } + { + let retrying_heals = self.retrying_heals.lock().await; + if retrying_heals + .values() + .any(|retrying| heal_type_matches_path(&retrying.request.heal_type, heal_path)) + { + return true; + } + } + let mut completed_heals = self.completed_heals.lock().await; prune_completed_heal_statuses(&mut completed_heals); completed_heals @@ -1040,8 +1378,9 @@ impl HealManager { } pub async fn get_task_progress(&self, task_id: &str) -> Result { + let canonical_task_id = self.canonical_task_id(task_id).await; let active_heals = self.active_heals.lock().await; - if let Some(task) = active_heals.get(task_id) { + if let Some(task) = active_heals.get(&canonical_task_id) { Ok(task.get_progress().await) } else { Err(Error::TaskNotFound { @@ -1052,37 +1391,62 @@ impl HealManager { /// Cancel task pub async fn cancel_task(&self, task_id: &str) -> Result<()> { + let canonical_task_id = self.canonical_task_id(task_id).await; { let mut active_heals = self.active_heals.lock().await; - if let Some(task) = active_heals.get(task_id) { + if let Some(task) = active_heals.get(&canonical_task_id) { task.cancel().await?; - active_heals.remove(task_id); + active_heals.remove(&canonical_task_id); publish_active_heal_count(&active_heals); info!( target: "rustfs::heal::manager", event = EVENT_HEAL_MANAGER_STATE, component = LOG_COMPONENT_HEAL, subsystem = LOG_SUBSYSTEM_MANAGER, - task_id, + task_id = %canonical_task_id, state = "cancelled_active_task", "Heal manager cancelled active task" ); + drop(active_heals); + self.remove_aliases_for_task(&canonical_task_id).await; + return Ok(()); + } + } + + { + let mut retrying_heals = self.retrying_heals.lock().await; + if let Some(retrying) = retrying_heals.remove(&canonical_task_id) { + retrying.cancel_token.cancel(); + drop(retrying_heals); + self.completed_heals.lock().await.remove(&canonical_task_id); + self.remove_aliases_for_task(&canonical_task_id).await; + info!( + target: "rustfs::heal::manager", + event = EVENT_HEAL_MANAGER_STATE, + component = LOG_COMPONENT_HEAL, + subsystem = LOG_SUBSYSTEM_MANAGER, + task_id = %canonical_task_id, + state = "cancelled_retrying_task", + "Heal manager cancelled retrying task" + ); return Ok(()); } } let mut queue = self.heal_queue.lock().await; - if queue.remove_request_id(task_id).is_some() { + if queue.remove_request_id(&canonical_task_id).is_some() { publish_heal_queue_length(&queue); info!( target: "rustfs::heal::manager", event = EVENT_HEAL_MANAGER_STATE, component = LOG_COMPONENT_HEAL, subsystem = LOG_SUBSYSTEM_MANAGER, - task_id, + task_id = %canonical_task_id, state = "cancelled_queued_task", "Heal manager cancelled queued task" ); + drop(queue); + self.remove_aliases_for_task(&canonical_task_id).await; return Ok(()); } @@ -1102,24 +1466,62 @@ impl HealManager { .filter_map(|(task_id, task)| heal_type_matches_path(&task.heal_type, heal_path).then_some(task_id.clone())) .collect::>(); - for task_id in task_ids { - if let Some(task) = active_heals.get(&task_id) { + for task_id in &task_ids { + if let Some(task) = active_heals.get(task_id) { task.cancel().await?; } - active_heals.remove(&task_id); + active_heals.remove(task_id); cancelled += 1; } if cancelled > 0 { publish_active_heal_count(&active_heals); } + drop(active_heals); + for task_id in &task_ids { + self.remove_aliases_for_task(task_id).await; + } + } + + { + let mut retrying_heals = self.retrying_heals.lock().await; + let task_ids = retrying_heals + .iter() + .filter_map(|(task_id, retrying)| { + heal_type_matches_path(&retrying.request.heal_type, heal_path).then_some(task_id.clone()) + }) + .collect::>(); + + for task_id in &task_ids { + if let Some(retrying) = retrying_heals.remove(task_id) { + retrying.cancel_token.cancel(); + cancelled += 1; + } + } + drop(retrying_heals); + + if !task_ids.is_empty() { + let mut completed_heals = self.completed_heals.lock().await; + for task_id in &task_ids { + completed_heals.remove(task_id); + } + drop(completed_heals); + + for task_id in &task_ids { + self.remove_aliases_for_task(task_id).await; + } + } } let mut queue = self.heal_queue.lock().await; let queued_cancelled = queue.remove_matching(|request| heal_type_matches_path(&request.heal_type, heal_path)); - if queued_cancelled > 0 { + if !queued_cancelled.is_empty() { publish_heal_queue_length(&queue); - cancelled += queued_cancelled; + cancelled += queued_cancelled.len(); + } + drop(queue); + for request in &queued_cancelled { + self.remove_aliases_for_task(&request.id).await; } if cancelled == 0 { @@ -1196,6 +1598,7 @@ impl HealManager { let heal_queue = self.heal_queue.clone(); let active_heals = self.active_heals.clone(); let completed_heals = self.completed_heals.clone(); + let retrying_heals = self.retrying_heals.clone(); let cancel_token = self.cancel_token.clone(); let statistics = self.statistics.clone(); let storage = self.storage.clone(); @@ -1219,10 +1622,32 @@ impl HealManager { break; } _ = notify.notified(), if event_driven_scheduler_enable => { - Self::process_heal_queue(&heal_queue, &active_heals, &completed_heals, &config, &statistics, &storage, ¬ify).await; + Self::process_heal_queue(HealQueueContext { + heal_queue: &heal_queue, + active_heals: &active_heals, + completed_heals: &completed_heals, + retrying_heals: &retrying_heals, + config: &config, + statistics: &statistics, + storage: &storage, + notify: ¬ify, + cancel_token: &cancel_token, + }) + .await; } _ = interval.tick() => { - Self::process_heal_queue(&heal_queue, &active_heals, &completed_heals, &config, &statistics, &storage, ¬ify).await; + Self::process_heal_queue(HealQueueContext { + heal_queue: &heal_queue, + active_heals: &active_heals, + completed_heals: &completed_heals, + retrying_heals: &retrying_heals, + config: &config, + statistics: &statistics, + storage: &storage, + notify: ¬ify, + cancel_token: &cancel_token, + }) + .await; } } } @@ -1450,15 +1875,19 @@ impl HealManager { /// Process heal queue /// Processes multiple tasks per cycle when capacity allows and queue has high-priority items - async fn process_heal_queue( - heal_queue: &Arc>, - active_heals: &Arc>>>, - completed_heals: &Arc>>, - config: &Arc>, - statistics: &Arc>, - storage: &Arc, - notify: &Arc, - ) { + async fn process_heal_queue(context: HealQueueContext<'_>) { + let HealQueueContext { + heal_queue, + active_heals, + completed_heals, + retrying_heals, + config, + statistics, + storage, + notify, + cancel_token, + } = context; + let config = config.read().await; let mut active_heals_guard = active_heals.lock().await; publish_active_heal_count(&active_heals_guard); @@ -1513,9 +1942,12 @@ impl HealManager { publish_active_heal_count(&active_heals_guard); update_task_running_metric_for_task(&active_heals_guard, task.as_ref()); let active_heals_clone = active_heals.clone(); + let heal_queue_clone = heal_queue.clone(); let completed_heals_clone = completed_heals.clone(); + let retrying_heals_clone = retrying_heals.clone(); let statistics_clone = statistics.clone(); let notify_clone = notify.clone(); + let manager_cancel_token = cancel_token.clone(); let task_type_label_for_spawn = task_type_label.clone(); let task_set_label_for_spawn = task_set_label.clone(); @@ -1534,7 +1966,8 @@ impl HealManager { "Heal scheduler task started" ); let result = task.execute().await; - match result { + let retry_request = retry_request_for_result(task.as_ref(), &result); + match &result { Ok(_) => { debug!( target: "rustfs::heal::manager", @@ -1549,25 +1982,51 @@ impl HealManager { ); } Err(e) => { - error!( - target: "rustfs::heal::manager", - event = EVENT_HEAL_SCHEDULER_STATE, - component = LOG_COMPONENT_HEAL, - subsystem = LOG_SUBSYSTEM_MANAGER, - task_id, - heal_type = %task_type_label_for_spawn, - set = %task_set_label_for_spawn, - state = "task_failed", - error = %e, - "Heal scheduler task failed" - ); + let will_retry = retry_request.is_some(); + if will_retry { + warn!( + target: "rustfs::heal::manager", + event = EVENT_HEAL_SCHEDULER_STATE, + component = LOG_COMPONENT_HEAL, + subsystem = LOG_SUBSYSTEM_MANAGER, + task_id, + heal_type = %task_type_label_for_spawn, + set = %task_set_label_for_spawn, + state = "task_retrying", + retry_attempt = task.retry_attempts.saturating_add(1), + error = %e, + "Heal scheduler task retrying" + ); + } else { + error!( + target: "rustfs::heal::manager", + event = EVENT_HEAL_SCHEDULER_STATE, + component = LOG_COMPONENT_HEAL, + subsystem = LOG_SUBSYSTEM_MANAGER, + task_id, + heal_type = %task_type_label_for_spawn, + set = %task_set_label_for_spawn, + state = "task_failed", + error = %e, + "Heal scheduler task failed" + ); + } } } + let retry_request_for_status = retry_request.as_ref().map(|(request, _, error)| HealTaskStatus::Retrying { + error: error.clone(), + retry_attempt: request.retry_attempts, + }); + let retry_request_for_queue = retry_request; let mut active_heals_guard = active_heals_clone.lock().await; if let Some(completed_task) = active_heals_guard.remove(&task_id) { publish_active_heal_count(&active_heals_guard); update_task_running_metric_for_task(&active_heals_guard, completed_task.as_ref()); - let completed_status = completed_task.get_status().await; + let completed_status = if let Some(status) = retry_request_for_status { + status + } else { + completed_task.get_status().await + }; let completed_status_entry = CompletedHealStatus { heal_type: completed_task.heal_type.clone(), status: completed_status.clone(), @@ -1583,12 +2042,118 @@ impl HealManager { HealTaskStatus::Completed => { stats.update_task_completion(true); } + HealTaskStatus::Retrying { .. } => {} _ => { stats.update_task_completion(false); } } stats.update_running_tasks(active_heals_guard.len() as u64); } + + if let Some((retry_request, retry_delay, retry_error)) = retry_request_for_queue { + let retry_request_id = retry_request.id.clone(); + let retry_attempt = retry_request.retry_attempts; + let retry_key = PriorityHealQueue::make_dedup_key(&retry_request); + let retry_priority = retry_request.priority; + let retry_active_heals = active_heals_clone.clone(); + let retry_heal_queue = heal_queue_clone.clone(); + let retrying_heals_for_spawn = retrying_heals_clone.clone(); + let retry_notify = notify_clone.clone(); + let retry_cancel_token = CancellationToken::new(); + let retry_manager_cancel_token = manager_cancel_token.clone(); + retrying_heals_clone.lock().await.insert( + retry_request_id.clone(), + RetryingHeal { + request: retry_request.clone(), + error: retry_error.clone(), + cancel_token: retry_cancel_token.clone(), + }, + ); + tokio::spawn(async move { + tokio::select! { + _ = retry_cancel_token.cancelled() => { + info!( + target: "rustfs::heal::manager", + event = EVENT_HEAL_QUEUE_ADMISSION, + component = LOG_COMPONENT_HEAL, + subsystem = LOG_SUBSYSTEM_MANAGER, + request_id = %retry_request_id, + priority = ?retry_priority, + retry_attempt, + result = "retry_cancelled", + "Heal retry admission decided" + ); + return; + } + _ = retry_manager_cancel_token.cancelled() => { + retrying_heals_for_spawn.lock().await.remove(&retry_request_id); + return; + } + _ = sleep(retry_delay) => {} + } + + { + let mut retrying_heals_guard = retrying_heals_for_spawn.lock().await; + if retrying_heals_guard.remove(&retry_request_id).is_none() { + return; + } + } + + { + let active_heals_guard = retry_active_heals.lock().await; + if active_heals_contains_dedup_key(&active_heals_guard, &retry_key) { + info!( + target: "rustfs::heal::manager", + event = EVENT_HEAL_QUEUE_ADMISSION, + component = LOG_COMPONENT_HEAL, + subsystem = LOG_SUBSYSTEM_MANAGER, + request_id = %retry_request_id, + priority = ?retry_priority, + retry_attempt, + result = "retry_merged_active_duplicate", + "Heal retry admission decided" + ); + return; + } + } + + let mut queue = retry_heal_queue.lock().await; + let outcome = queue.push(retry_request); + publish_heal_queue_length(&queue); + match outcome { + QueuePushOutcome::Accepted => { + info!( + target: "rustfs::heal::manager", + event = EVENT_HEAL_QUEUE_ADMISSION, + component = LOG_COMPONENT_HEAL, + subsystem = LOG_SUBSYSTEM_MANAGER, + request_id = %retry_request_id, + priority = ?retry_priority, + retry_attempt, + retry_delay_ms = retry_delay.as_millis(), + error = %retry_error, + result = "retry_enqueued", + "Heal retry admission decided" + ); + drop(queue); + retry_notify.notify_one(); + } + QueuePushOutcome::Merged => { + info!( + target: "rustfs::heal::manager", + event = EVENT_HEAL_QUEUE_ADMISSION, + component = LOG_COMPONENT_HEAL, + subsystem = LOG_SUBSYSTEM_MANAGER, + request_id = %retry_request_id, + priority = ?retry_priority, + retry_attempt, + result = "retry_merged_duplicate", + "Heal retry admission decided" + ); + } + } + }); + } notify_clone.notify_one(); }); tasks_started += 1; @@ -2244,6 +2809,246 @@ mod tests { ); } + #[tokio::test] + async fn test_submit_heal_request_returns_merged_for_active_duplicate() { + let storage: Arc = Arc::new(MockStorage); + let manager = HealManager::new(storage.clone(), None); + let active_request = HealRequest::object("bucket".to_string(), "object".to_string(), None); + let active_task = Arc::new(HealTask::from_request(active_request, storage)); + manager.active_heals.lock().await.insert(active_task.id.clone(), active_task); + + let duplicate_request = HealRequest::object("bucket".to_string(), "object".to_string(), None); + + assert_eq!( + manager + .submit_heal_request(duplicate_request) + .await + .expect("active duplicate should produce admission result"), + HealAdmissionResult::Merged + ); + assert_eq!(manager.get_queue_length().await, 0); + } + + #[tokio::test] + async fn test_active_duplicate_token_can_query_and_cancel_original_task() { + let storage: Arc = Arc::new(MockStorage); + let manager = HealManager::new(storage.clone(), None); + let active_request = HealRequest::object("bucket".to_string(), "object".to_string(), None); + let active_task = Arc::new(HealTask::from_request(active_request, storage)); + let active_task_id = active_task.id.clone(); + manager.active_heals.lock().await.insert(active_task_id.clone(), active_task); + + let duplicate_request = HealRequest::object("bucket".to_string(), "object".to_string(), None); + let duplicate_task_id = duplicate_request.id.clone(); + + assert_eq!( + manager + .submit_heal_request(duplicate_request) + .await + .expect("active duplicate should produce admission result"), + HealAdmissionResult::Merged + ); + assert_eq!( + manager + .get_task_status_for_path("bucket/object", &duplicate_task_id) + .await + .expect("duplicate token should query merged active task"), + HealTaskStatus::Pending + ); + + manager + .cancel_task(&duplicate_task_id) + .await + .expect("duplicate token should cancel merged active task"); + + assert!(manager.active_heals.lock().await.get(&active_task_id).is_none()); + assert!(matches!(manager.get_task_status(&active_task_id).await, Err(Error::TaskNotFound { .. }))); + } + + #[tokio::test] + async fn test_queued_duplicate_token_can_query_and_cancel_original_request() { + let storage: Arc = Arc::new(MockStorage); + let manager = HealManager::new(storage, None); + let original_request = HealRequest::object("bucket".to_string(), "object".to_string(), None); + let original_task_id = original_request.id.clone(); + let duplicate_request = HealRequest::object("bucket".to_string(), "object".to_string(), None); + let duplicate_task_id = duplicate_request.id.clone(); + + assert_eq!( + manager + .submit_heal_request(original_request) + .await + .expect("original request should be accepted"), + HealAdmissionResult::Accepted + ); + assert_eq!( + manager + .submit_heal_request(duplicate_request) + .await + .expect("queued duplicate should produce admission result"), + HealAdmissionResult::Merged + ); + assert_eq!( + manager + .get_task_status_for_path("bucket/object", &duplicate_task_id) + .await + .expect("duplicate token should query merged queued task"), + HealTaskStatus::Pending + ); + + manager + .cancel_task(&duplicate_task_id) + .await + .expect("duplicate token should cancel merged queued request"); + + assert!(matches!( + manager.get_task_status(&original_task_id).await, + Err(Error::TaskNotFound { .. }) + )); + } + + #[test] + fn test_retry_request_for_recoverable_lock_timeout() { + let storage: Arc = Arc::new(MockStorage); + let task = HealTask::from_request(HealRequest::object("bucket".to_string(), "object".to_string(), None), storage); + let result = Err(Error::TaskExecutionFailed { + message: "Failed to heal object bucket/object: Lock acquisition timeout".to_string(), + }); + + let (retry_request, retry_delay, retry_error) = + retry_request_for_result(&task, &result).expect("lock timeout should be retryable"); + + assert_eq!(retry_request.id, task.id); + assert_eq!(retry_request.retry_attempts, 1); + assert_eq!(retry_request.priority, task.priority); + assert!(retry_delay > Duration::ZERO); + assert!(retry_error.contains("Lock acquisition timeout")); + } + + #[test] + fn test_retry_request_for_recoverable_error_stops_at_limit() { + let storage: Arc = Arc::new(MockStorage); + let mut request = HealRequest::object("bucket".to_string(), "object".to_string(), None); + request.retry_attempts = MAX_RECOVERABLE_HEAL_RETRIES; + let task = HealTask::from_request(request, storage); + let result = Err(Error::TaskExecutionFailed { + message: "Remote lock RPC timed out".to_string(), + }); + + assert!(retry_request_for_result(&task, &result).is_none()); + } + + async fn insert_retrying_request(manager: &HealManager, request: HealRequest) -> CancellationToken { + let task_id = request.id.clone(); + let cancel_token = CancellationToken::new(); + manager.retrying_heals.lock().await.insert( + task_id.clone(), + RetryingHeal { + request: request.clone(), + error: "Lock acquisition timeout".to_string(), + cancel_token: cancel_token.clone(), + }, + ); + manager.completed_heals.lock().await.insert( + task_id, + CompletedHealStatus { + heal_type: request.heal_type, + status: HealTaskStatus::Retrying { + error: "Lock acquisition timeout".to_string(), + retry_attempt: request.retry_attempts, + }, + result_items: Vec::new(), + completed_at: SystemTime::now(), + }, + ); + cancel_token + } + + #[tokio::test] + async fn test_cancel_task_cancels_retrying_backoff() { + let storage: Arc = Arc::new(MockStorage); + let manager = HealManager::new(storage, None); + let mut request = HealRequest::bucket("bucket".to_string()); + request.retry_attempts = 1; + let task_id = request.id.clone(); + let cancel_token = insert_retrying_request(&manager, request).await; + + assert!(matches!( + manager + .get_task_status(&task_id) + .await + .expect("retrying task should be queryable"), + HealTaskStatus::Retrying { .. } + )); + + manager + .cancel_task(&task_id) + .await + .expect("retrying task should be cancellable by token"); + + assert!(cancel_token.is_cancelled()); + assert!(manager.retrying_heals.lock().await.get(&task_id).is_none()); + assert!(matches!(manager.get_task_status(&task_id).await, Err(Error::TaskNotFound { .. }))); + } + + #[tokio::test] + async fn test_cancel_tasks_for_path_cancels_retrying_backoff() { + let storage: Arc = Arc::new(MockStorage); + let manager = HealManager::new(storage, None); + let mut request = HealRequest::bucket("bucket".to_string()); + request.retry_attempts = 1; + let task_id = request.id.clone(); + let cancel_token = insert_retrying_request(&manager, request).await; + + assert_eq!( + manager + .cancel_tasks_for_path("bucket") + .await + .expect("retrying task should be cancellable by path"), + 1 + ); + + assert!(cancel_token.is_cancelled()); + assert!(manager.retrying_heals.lock().await.get(&task_id).is_none()); + assert!(matches!(manager.get_task_status(&task_id).await, Err(Error::TaskNotFound { .. }))); + } + + #[tokio::test] + async fn test_retrying_duplicate_token_can_query_and_cancel_original_retry() { + let storage: Arc = Arc::new(MockStorage); + let manager = HealManager::new(storage, None); + let mut original_request = HealRequest::bucket("bucket".to_string()); + original_request.retry_attempts = 1; + let original_task_id = original_request.id.clone(); + let cancel_token = insert_retrying_request(&manager, original_request).await; + + let duplicate_request = HealRequest::bucket("bucket".to_string()); + let duplicate_task_id = duplicate_request.id.clone(); + + assert_eq!( + manager + .submit_heal_request(duplicate_request) + .await + .expect("retrying duplicate should produce admission result"), + HealAdmissionResult::Merged + ); + assert!(matches!( + manager + .get_task_status_for_path("bucket", &duplicate_task_id) + .await + .expect("duplicate token should query merged retrying task"), + HealTaskStatus::Retrying { .. } + )); + + manager + .cancel_task(&duplicate_task_id) + .await + .expect("duplicate token should cancel merged retrying task"); + + assert!(cancel_token.is_cancelled()); + assert!(manager.retrying_heals.lock().await.get(&original_task_id).is_none()); + } + #[tokio::test] async fn test_get_task_status_reports_pending_for_queued_request() { let storage: Arc = Arc::new(MockStorage); diff --git a/crates/heal/src/heal/task.rs b/crates/heal/src/heal/task.rs index 7e7f8ca71..2b2cd7b65 100644 --- a/crates/heal/src/heal/task.rs +++ b/crates/heal/src/heal/task.rs @@ -148,6 +148,8 @@ pub enum HealTaskStatus { Pending, /// Running Running, + /// Retrying after a recoverable failure + Retrying { error: String, retry_attempt: u32 }, /// Completed Completed, /// Failed @@ -173,6 +175,8 @@ pub struct HealRequest { pub source: HealRequestSource, /// Whether this request should bypass queue admission dedup/full policies. pub force_start: bool, + /// Number of recoverable retry attempts already scheduled for this request. + pub retry_attempts: u32, /// Created time pub created_at: SystemTime, /// Queue admission time used for scheduler delay metrics @@ -189,6 +193,7 @@ impl HealRequest { priority, source: HealRequestSource::Internal, force_start: false, + retry_attempts: 0, created_at: now, enqueued_at: now, } @@ -239,6 +244,8 @@ pub struct HealTask { pub priority: HealPriority, /// Origin inherited from the request pub source: HealRequestSource, + /// Number of recoverable retry attempts already scheduled for this task. + pub retry_attempts: u32, /// Task status pub status: Arc>, /// Progress tracking @@ -269,6 +276,7 @@ impl HealTask { options: request.options, priority: request.priority, source: request.source, + retry_attempts: request.retry_attempts, status: Arc::new(RwLock::new(HealTaskStatus::Pending)), progress: Arc::new(RwLock::new(HealProgress::new())), result_items: Arc::new(RwLock::new(Vec::new())), @@ -282,6 +290,20 @@ impl HealTask { } } + pub fn retry_request(&self) -> HealRequest { + HealRequest { + id: self.id.clone(), + heal_type: self.heal_type.clone(), + options: self.options.clone(), + priority: self.priority, + source: self.source, + force_start: false, + retry_attempts: self.retry_attempts.saturating_add(1), + created_at: self.created_at, + enqueued_at: SystemTime::now(), + } + } + pub fn metric_type_label(&self) -> &'static str { match &self.heal_type { HealType::Object { .. } => "object", diff --git a/rustfs/src/admin/handlers/heal.rs b/rustfs/src/admin/handlers/heal.rs index c47954845..7b48d2fd5 100644 --- a/rustfs/src/admin/handlers/heal.rs +++ b/rustfs/src/admin/handlers/heal.rs @@ -247,6 +247,11 @@ fn encode_heal_task_status( } fn build_heal_channel_request(hip: &HealInitParams) -> HealChannelRequest { + let recursive = if !hip.bucket.is_empty() && hip.obj_prefix.is_empty() { + true + } else { + hip.hs.recursive + }; let mut heal_request = rustfs_common::heal_channel::create_heal_request( hip.bucket.clone(), if hip.obj_prefix.is_empty() { @@ -264,7 +269,7 @@ fn build_heal_channel_request(hip: &HealInitParams) -> HealChannelRequest { heal_request.remove_corrupted = Some(hip.hs.remove); heal_request.recreate_missing = Some(hip.hs.recreate); heal_request.update_parity = Some(hip.hs.update_parity); - heal_request.recursive = Some(hip.hs.recursive); + heal_request.recursive = Some(recursive); heal_request.dry_run = Some(hip.hs.dry_run); heal_request.source = HealRequestSource::Admin; heal_request @@ -831,6 +836,27 @@ mod tests { assert_eq!(request.set_index, Some(2)); } + #[test] + fn test_bucket_heal_channel_request_defaults_to_recursive() { + let hip = HealInitParams { + bucket: "bucket".to_string(), + obj_prefix: String::new(), + hs: HealOpts { + recursive: false, + scan_mode: HealScanMode::Deep, + recreate: true, + ..Default::default() + }, + ..Default::default() + }; + + let request = build_heal_channel_request(&hip); + + assert_eq!(request.bucket, "bucket"); + assert_eq!(request.object_prefix, None); + assert_eq!(request.recursive, Some(true)); + } + #[test] fn test_extract_heal_init_params_allows_root_heal_target() { let uri: Uri = "/rustfs/admin/v3/heal/".parse().expect("uri should parse");