From 1a6b870eb5df014a1806b41d41b677be829762c4 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E9=A9=AC=E7=99=BB=E5=B1=B1?= Date: Sat, 22 Aug 2026 02:39:36 +0800 Subject: [PATCH] fix(heal): supervise scheduler task panics --- crates/heal/src/heal/manager/scheduler.rs | 1024 +++++++++++++++------ crates/heal/src/heal/manager/tests.rs | 230 ++++- 2 files changed, 948 insertions(+), 306 deletions(-) diff --git a/crates/heal/src/heal/manager/scheduler.rs b/crates/heal/src/heal/manager/scheduler.rs index fbee9f212..5ada4f29b 100644 --- a/crates/heal/src/heal/manager/scheduler.rs +++ b/crates/heal/src/heal/manager/scheduler.rs @@ -13,6 +13,341 @@ // limitations under the License. /// The heal scheduler: queue consumption loop and its skip/metric helpers. use super::*; +use futures::FutureExt; +use std::panic::AssertUnwindSafe; + +pub(super) const PANICKED_HEAL_TASK_ERROR: &str = "heal task panicked"; + +#[cfg(test)] +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub(super) enum SchedulerPanicPoint { + RetryChild, + Cleanup, +} + +#[cfg(test)] +#[derive(Clone, Debug, PartialEq, Eq)] +struct SchedulerPanicHook { + point: SchedulerPanicPoint, + task_id: String, +} + +#[cfg(test)] +static SCHEDULER_PANIC_POINT: LazyLock>> = LazyLock::new(|| StdMutex::new(None)); + +#[cfg(test)] +pub(super) fn arm_scheduler_panic(point: SchedulerPanicPoint, task_id: &str) { + *SCHEDULER_PANIC_POINT + .lock() + .expect("scheduler panic hook lock should not be poisoned") = Some(SchedulerPanicHook { + point, + task_id: task_id.to_string(), + }); +} + +#[cfg(test)] +pub(super) fn clear_scheduler_panic() { + *SCHEDULER_PANIC_POINT + .lock() + .expect("scheduler panic hook lock should not be poisoned") = None; +} + +#[cfg(test)] +fn panic_if_armed(point: SchedulerPanicPoint, task_id: &str) { + let mut hook = SCHEDULER_PANIC_POINT + .lock() + .expect("scheduler panic hook lock should not be poisoned"); + if hook + .as_ref() + .is_some_and(|hook| hook.point == point && hook.task_id == task_id) + { + hook.take(); + drop(hook); + panic!("test-only scheduler panic hook"); + } +} + +#[derive(Clone)] +pub(super) struct PanicCleanupState { + pub(super) active_heals: Arc>>>, + pub(super) heal_queue: Arc>, + pub(super) completed_heals: Arc>>>, + pub(super) task_aliases: Arc>>, + pub(super) retrying_heals: Arc>>, + pub(super) mrf_repair_notice_targets: Arc>>>, + pub(super) replacement_recovery_anchors: Arc>>, + pub(super) statistics: Arc>, +} + +pub(super) async fn finish_panicked_heal_task(task: Arc, task_id: String, state: PanicCleanupState) { + // A panic can interrupt any point between the active/retrying/completed + // handoff. Each operation is intentionally idempotent so an outer unwind + // handler can safely repair a partially completed handoff. + let cancelled = task.cancel_token.is_cancelled(); + task.cancel_token.cancel(); + let current_status = AssertUnwindSafe(task.get_status()).catch_unwind().await.ok(); + let mut terminal_status = match current_status { + Some(status) + if matches!( + status, + HealTaskStatus::Completed | HealTaskStatus::Failed { .. } | HealTaskStatus::Cancelled | HealTaskStatus::Timeout + ) => + { + status + } + _ if cancelled => HealTaskStatus::Cancelled, + _ => HealTaskStatus::Failed { + error: PANICKED_HEAL_TASK_ERROR.to_string(), + }, + }; + let _ = AssertUnwindSafe(async { + *task.completed_at.write().await = Some(SystemTime::now()); + }) + .catch_unwind() + .await; + let _ = AssertUnwindSafe(async { + *task.status.write().await = terminal_status.clone(); + }) + .catch_unwind() + .await; + + let (removed_active, active_count) = AssertUnwindSafe(async { + let mut active = state.active_heals.lock().await; + let removed_task = active.remove(&task_id); + if let Some(removed_task) = removed_task.as_ref() { + update_task_running_metric_for_task(&active, removed_task.as_ref()); + } + let removed = removed_task.is_some(); + publish_active_heal_count(&active); + (removed, active.len()) + }) + .catch_unwind() + .await + .unwrap_or((false, 0)); + + let _ = AssertUnwindSafe(async { + if let Some(retrying) = state.retrying_heals.lock().await.remove(&task_id) { + retrying.cancel_token.cancel(); + } + }) + .catch_unwind() + .await; + + let _ = AssertUnwindSafe(async { + let mut queue = state.heal_queue.lock().await; + queue.remove_request_id(&task_id); + publish_heal_queue_length(&queue); + }) + .catch_unwind() + .await; + let _ = AssertUnwindSafe(async { + state + .replacement_recovery_anchors + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()) + .remove(&task_id); + }) + .catch_unwind() + .await; + let _ = AssertUnwindSafe(async { + state + .task_aliases + .lock() + .await + .retain(|alias_id, alias| alias_id != &task_id && alias.task_id != task_id); + }) + .catch_unwind() + .await; + + // An explicit cancellation can race with the supervisor after the first + // status snapshot. Let the cancellation win if it has already published + // its terminal state before we archive the panic. + if let Ok(status) = AssertUnwindSafe(task.get_status()).catch_unwind().await + && matches!(status, HealTaskStatus::Cancelled) + { + terminal_status = HealTaskStatus::Cancelled; + } + + let _ = AssertUnwindSafe(async { + if matches!(&terminal_status, HealTaskStatus::Completed) { + emit_mrf_repaired_events(take_mrf_repair_notice_targets(&state.mrf_repair_notice_targets, &task_id)); + } else { + lock_mrf_repair_notice_targets(&state.mrf_repair_notice_targets).remove(&task_id); + } + }) + .catch_unwind() + .await; + + let successful = matches!(terminal_status, HealTaskStatus::Completed); + let completed_status = AssertUnwindSafe(async { + let seqed_items = task.get_seqed_result_items().await; + let (next_seq, min_seq) = task.result_seq_cursors(); + CompletedHealStatus { + heal_type: task.heal_type.clone(), + status: terminal_status.clone(), + result_items_truncated: task.result_items_truncated(), + completed_at: SystemTime::now(), + seqed_items, + next_seq, + min_seq, + } + }) + .catch_unwind() + .await + .unwrap_or_else(|_| CompletedHealStatus { + heal_type: task.heal_type.clone(), + status: terminal_status.clone(), + result_items_truncated: false, + completed_at: SystemTime::now(), + seqed_items: Vec::new(), + next_seq: 0, + min_seq: 0, + }); + let archived = AssertUnwindSafe(async { + let mut completed = state.completed_heals.lock().await; + // cancel_task removes active work and preserves its historical + // TaskNotFound behavior. If that removal already won, the panic + // supervisor must only clear stale secondary state. + if matches!(&terminal_status, HealTaskStatus::Cancelled) && !removed_active { + return false; + } + // The retained terminal entry is the finish-once CAS shared by the + // worker and both panic supervisors; never replace a terminal result. + let replace_existing = completed + .get(&task_id) + .map(|existing| { + !matches!( + existing.status, + HealTaskStatus::Completed + | HealTaskStatus::Failed { .. } + | HealTaskStatus::Cancelled + | HealTaskStatus::Timeout + ) + }) + .unwrap_or(true); + if replace_existing || !completed.contains_key(&task_id) { + prune_completed_heal_statuses(&mut completed); + completed.insert(task_id.clone(), Arc::new(completed_status)); + true + } else { + false + } + }) + .catch_unwind() + .await + .unwrap_or(false); + + let completed_progress = AssertUnwindSafe(task.get_progress()).catch_unwind().await.unwrap_or_default(); + if archived { + let _ = AssertUnwindSafe(async { + let mut stats = state.statistics.write().await; + if successful { + stats.update_task_completion(true); + stats.add_healed_objects(completed_progress.objects_healed, completed_progress.bytes_processed); + } else { + stats.update_task_completion(false); + } + stats.update_running_tasks(usize_to_u64_saturated(active_count)); + }) + .catch_unwind() + .await; + } +} + +pub(super) async fn finish_panicked_retry_child( + retry_request_id: String, + heal_type: HealType, + retry_cancel_token: CancellationToken, + state: PanicCleanupState, +) { + let _ = AssertUnwindSafe(async { + state.retrying_heals.lock().await.remove(&retry_request_id); + }) + .catch_unwind() + .await; + let _ = AssertUnwindSafe(async { + let mut queue = state.heal_queue.lock().await; + queue.remove_request_id(&retry_request_id); + publish_heal_queue_length(&queue); + }) + .catch_unwind() + .await; + let archived = AssertUnwindSafe(async { + let mut completed = state.completed_heals.lock().await; + // Recheck after acquiring the completed lock: cancel_task cancels the + // retry token before waiting on this lock, so a late panic must not + // recreate a terminal Failed entry after cancellation wins. + if retry_cancel_token.is_cancelled() { + return false; + } + if let Some(existing) = completed.get(&retry_request_id) { + let mut updated = (**existing).clone(); + if !matches!( + updated.status, + HealTaskStatus::Completed | HealTaskStatus::Failed { .. } | HealTaskStatus::Cancelled | HealTaskStatus::Timeout + ) { + updated.status = HealTaskStatus::Failed { + error: PANICKED_HEAL_TASK_ERROR.to_string(), + }; + updated.completed_at = SystemTime::now(); + completed.insert(retry_request_id.clone(), Arc::new(updated)); + true + } else { + false + } + } else { + prune_completed_heal_statuses(&mut completed); + completed.insert( + retry_request_id.clone(), + Arc::new(CompletedHealStatus { + heal_type, + status: HealTaskStatus::Failed { + error: PANICKED_HEAL_TASK_ERROR.to_string(), + }, + result_items_truncated: false, + completed_at: SystemTime::now(), + seqed_items: Vec::new(), + next_seq: 0, + min_seq: 0, + }), + ); + true + } + }) + .catch_unwind() + .await + .unwrap_or(false); + let _ = AssertUnwindSafe(async { + state + .task_aliases + .lock() + .await + .retain(|alias_id, alias| alias_id != &retry_request_id && alias.task_id != retry_request_id); + }) + .catch_unwind() + .await; + let _ = AssertUnwindSafe(async { + lock_mrf_repair_notice_targets(&state.mrf_repair_notice_targets).remove(&retry_request_id); + }) + .catch_unwind() + .await; + let _ = AssertUnwindSafe(async { + state + .replacement_recovery_anchors + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()) + .remove(&retry_request_id); + }) + .catch_unwind() + .await; + if archived { + let _ = AssertUnwindSafe(async { + state.statistics.write().await.update_task_completion(false); + }) + .catch_unwind() + .await; + } +} impl HealManager { /// Start scheduler @@ -193,54 +528,39 @@ impl HealManager { let task_type_label_for_spawn = task_type_label.clone(); let task_set_label_for_spawn = task_set_label.clone(); let config_for_spawn = config.clone(); + let panic_task = task.clone(); + let panic_task_id = task_id.clone(); + let panic_state = PanicCleanupState { + active_heals: active_heals.clone(), + heal_queue: heal_queue.clone(), + completed_heals: completed_heals.clone(), + task_aliases: task_aliases.clone(), + retrying_heals: retrying_heals.clone(), + mrf_repair_notice_targets: mrf_repair_notice_targets.clone(), + replacement_recovery_anchors: replacement_recovery_anchors.clone(), + statistics: statistics.clone(), + }; // start heal task tokio::spawn(async move { - debug!( - target: "rustfs::heal::manager", - event = EVENT_HEAL_SCHEDULER_STATE, - component = LOG_COMPONENT_HEAL, - subsystem = LOG_SUBSYSTEM_MANAGER, - task_id, - priority = ?task_priority, - heal_type = %task_type_label_for_spawn, - set = %task_set_label_for_spawn, - state = "task_started", - "Heal scheduler task started" - ); - let result = task.execute().await; - let retry_request = retry_request_for_result_with_budget(task.as_ref(), &result).await; - match &result { - Ok(_) => { - debug!( - 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_completed", - "Heal scheduler task completed" - ); - } - Err(e) => { - let will_retry = retry_request.is_some(); - if will_retry { - demote_to_debug_when!(task.heal_type.is_per_object(), 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!( + let scheduler_task = async move { + debug!( + target: "rustfs::heal::manager", + event = EVENT_HEAL_SCHEDULER_STATE, + component = LOG_COMPONENT_HEAL, + subsystem = LOG_SUBSYSTEM_MANAGER, + task_id, + priority = ?task_priority, + heal_type = %task_type_label_for_spawn, + set = %task_set_label_for_spawn, + state = "task_started", + "Heal scheduler task started" + ); + let result = task.execute().await; + let retry_request = retry_request_for_result_with_budget(task.as_ref(), &result).await; + match &result { + Ok(_) => { + debug!( target: "rustfs::heal::manager", event = EVENT_HEAL_SCHEDULER_STATE, component = LOG_COMPONENT_HEAL, @@ -248,284 +568,378 @@ impl HealManager { task_id, heal_type = %task_type_label_for_spawn, set = %task_set_label_for_spawn, - state = "task_failed", - error = %e, - "Heal scheduler task failed" + state = "task_completed", + "Heal scheduler task completed" ); } - } - } - 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 retry_cancel_token = retry_request_for_queue.as_ref().map(|_| CancellationToken::new()); - if retry_request_for_queue.is_none() { - replacement_recovery_anchors_clone - .lock() - .unwrap_or_else(|poisoned| poisoned.into_inner()) - .remove(&task_id); - } - let mut active_heals_guard = active_heals_clone.lock().await; - // Keep retry ownership continuous: status snapshots acquire - // these locks in the same active -> retrying order. - let mut retrying_heals_guard = if let (Some((request, _, error)), Some(cancel_token)) = - (retry_request_for_queue.as_ref(), retry_cancel_token.as_ref()) - { - let mut retrying = retrying_heals_clone.lock().await; - if active_heals_guard.contains_key(&task_id) { - retrying.insert( - request.id.clone(), - RetryingHeal { - request: request.clone(), - error: error.clone(), - cancel_token: cancel_token.clone(), - }, - ); - #[cfg(test)] - pause_retry_ownership_transition(&task_id, false).await; - } - Some(retrying) - } else { - None - }; - let completed_task = active_heals_guard.remove(&task_id); - if let Some(completed_task) = completed_task.as_ref() { - publish_active_heal_count(&active_heals_guard); - update_task_running_metric_for_task(&active_heals_guard, completed_task.as_ref()); - } - let active_count = active_heals_guard.len(); - drop(retrying_heals_guard.take()); - drop(active_heals_guard); - - if let Some(completed_task) = completed_task { - let completed_status = if let Some(status) = retry_request_for_status { - status - } else { - completed_task.get_status().await - }; - let terminal_completion = !matches!(completed_status, HealTaskStatus::Retrying { .. }); - let successful_completion = matches!(completed_status, HealTaskStatus::Completed); - let completed_progress = completed_task.get_progress().await; - // Single snapshot of the retained window: the task is - // finished and already off the active map, so there is - // no concurrent writer to race with. - let seqed_items = completed_task.get_seqed_result_items().await; - let (next_seq, min_seq) = completed_task.result_seq_cursors(); - let completed_status_entry = CompletedHealStatus { - heal_type: completed_task.heal_type.clone(), - status: completed_status.clone(), - result_items_truncated: completed_task.result_items_truncated(), - completed_at: SystemTime::now(), - seqed_items, - next_seq, - min_seq, - }; - let mut completed_heals_guard = completed_heals_clone.lock().await; - prune_completed_heal_statuses(&mut completed_heals_guard); - completed_heals_guard.insert(task_id.clone(), Arc::new(completed_status_entry)); - drop(completed_heals_guard); - // update statistics - let mut stats = statistics_clone.write().await; - match completed_status { - HealTaskStatus::Completed => { - stats.update_task_completion(true); - stats.add_healed_objects(completed_progress.objects_healed, completed_progress.bytes_processed); - } - HealTaskStatus::Retrying { .. } => {} - _ => { - stats.update_task_completion(false); - } - } - stats.update_running_tasks(usize_to_u64_saturated(active_count)); - drop(stats); - if terminal_completion { - let notice_targets = take_mrf_repair_notice_targets(&mrf_repair_notice_targets_clone, &task_id); - if successful_completion { - emit_mrf_repaired_events(notice_targets); - } - task_aliases_clone - .lock() - .await - .retain(|alias_id, alias| alias_id != &task_id && alias.task_id != task_id); - } - } - - if let (Some((retry_request, retry_delay, retry_error)), Some(retry_cancel_token)) = - (retry_request_for_queue, retry_cancel_token) - { - 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_task_aliases = task_aliases_clone.clone(); - let retry_mrf_repair_notice_targets = mrf_repair_notice_targets_clone.clone(); - let retry_completed_heals = completed_heals_clone.clone(); - let retry_notify = notify_clone.clone(); - let retry_manager_cancel_token = manager_cancel_token.clone(); - let retry_config = config_for_spawn.clone(); - tokio::spawn(async move { - loop { - tokio::select! { - _ = retry_cancel_token.cancelled() => { - debug!( - 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 retrying_heals_guard = retrying_heals_for_spawn.lock().await; - if !retrying_heals_guard.contains_key(&retry_request_id) { - return; - } - } - - let active_duplicate_task_id = { - let active_heals_guard = retry_active_heals.lock().await; - active_heal_for_dedup_key(&active_heals_guard, &retry_key).map(|(task_id, _)| task_id) - }; - if let Some(active_duplicate_task_id) = active_duplicate_task_id { - retrying_heals_for_spawn.lock().await.remove(&retry_request_id); - move_mrf_repair_notice_targets( - &retry_mrf_repair_notice_targets, - &retry_request_id, - &active_duplicate_task_id, - ); - debug!( - target: "rustfs::heal::manager", - event = EVENT_HEAL_QUEUE_ADMISSION, + Err(e) => { + let will_retry = retry_request.is_some(); + if will_retry { + demote_to_debug_when!(task.heal_type.is_per_object(), warn, target: "rustfs::heal::manager", { + event = EVENT_HEAL_SCHEDULER_STATE, 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" + 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" ); - return; } + } + } + 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 retry_cancel_token = retry_request_for_queue.as_ref().map(|_| CancellationToken::new()); + if retry_request_for_queue.is_none() { + replacement_recovery_anchors_clone + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()) + .remove(&task_id); + } + let mut active_heals_guard = active_heals_clone.lock().await; + // Keep retry ownership continuous: status snapshots acquire + // these locks in the same active -> retrying order. + let mut retrying_heals_guard = if let (Some((request, _, error)), Some(cancel_token)) = + (retry_request_for_queue.as_ref(), retry_cancel_token.as_ref()) + { + let mut retrying = retrying_heals_clone.lock().await; + if active_heals_guard.contains_key(&task_id) { + retrying.insert( + request.id.clone(), + RetryingHeal { + request: request.clone(), + error: error.clone(), + cancel_token: cancel_token.clone(), + }, + ); + #[cfg(test)] + pause_retry_ownership_transition(&task_id, false).await; + } + Some(retrying) + } else { + None + }; + let completed_task = active_heals_guard.remove(&task_id); + if let Some(completed_task) = completed_task.as_ref() { + publish_active_heal_count(&active_heals_guard); + update_task_running_metric_for_task(&active_heals_guard, completed_task.as_ref()); + } + let active_count = active_heals_guard.len(); + drop(retrying_heals_guard.take()); + drop(active_heals_guard); - let mut queue = retry_heal_queue.lock().await; - let admission_decision = - Self::admit_request_to_queue(&mut queue, retry_request.clone(), &retry_config, "retry"); - let admission = admission_decision.result; - let should_notify = matches!(admission, HealAdmissionResult::Accepted) - && retry_config.event_driven_scheduler_enable; - match admission { - HealAdmissionResult::Accepted => { - // Transfer ownership while holding queue -> retrying, - // matching operations_snapshot's lock order. - #[cfg(test)] - pause_retry_ownership_transition(&retry_request_id, true).await; - retrying_heals_for_spawn.lock().await.remove(&retry_request_id); - let displaced_task_id = admission_decision.displaced_task_id; - drop(queue); - if let Some(displaced_task_id) = displaced_task_id { - remove_task_aliases_for_task(&retry_task_aliases, &displaced_task_id).await; - remove_mrf_repair_notice_targets( - &retry_mrf_repair_notice_targets, - &displaced_task_id, - ); + if let Some(completed_task) = completed_task { + let completed_status = if let Some(status) = retry_request_for_status { + status + } else { + completed_task.get_status().await + }; + let terminal_completion = !matches!(completed_status, HealTaskStatus::Retrying { .. }); + let successful_completion = matches!(completed_status, HealTaskStatus::Completed); + let completed_progress = completed_task.get_progress().await; + // Single snapshot of the retained window: the task is + // finished and already off the active map, so there is + // no concurrent writer to race with. + let seqed_items = completed_task.get_seqed_result_items().await; + let (next_seq, min_seq) = completed_task.result_seq_cursors(); + let completed_status_entry = CompletedHealStatus { + heal_type: completed_task.heal_type.clone(), + status: completed_status.clone(), + result_items_truncated: completed_task.result_items_truncated(), + completed_at: SystemTime::now(), + seqed_items, + next_seq, + min_seq, + }; + let mut completed_heals_guard = completed_heals_clone.lock().await; + prune_completed_heal_statuses(&mut completed_heals_guard); + completed_heals_guard.insert(task_id.clone(), Arc::new(completed_status_entry)); + drop(completed_heals_guard); + // update statistics + let mut stats = statistics_clone.write().await; + match completed_status { + HealTaskStatus::Completed => { + stats.update_task_completion(true); + stats.add_healed_objects( + completed_progress.objects_healed, + completed_progress.bytes_processed, + ); + } + HealTaskStatus::Retrying { .. } => {} + _ => { + stats.update_task_completion(false); + } + } + stats.update_running_tasks(usize_to_u64_saturated(active_count)); + drop(stats); + #[cfg(test)] + panic_if_armed(SchedulerPanicPoint::Cleanup, &task_id); + if terminal_completion { + let notice_targets = take_mrf_repair_notice_targets(&mrf_repair_notice_targets_clone, &task_id); + if successful_completion { + emit_mrf_repaired_events(notice_targets); + } + task_aliases_clone + .lock() + .await + .retain(|alias_id, alias| alias_id != &task_id && alias.task_id != task_id); + } + } + + if let (Some((retry_request, retry_delay, retry_error)), Some(retry_cancel_token)) = + (retry_request_for_queue, retry_cancel_token) + { + 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_panic_heal_type = retry_request.heal_type.clone(); + let retry_panic_set_label = heal_request_set_metric_label(&retry_request); + 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_task_aliases = task_aliases_clone.clone(); + let retry_mrf_repair_notice_targets = mrf_repair_notice_targets_clone.clone(); + let retry_completed_heals = completed_heals_clone.clone(); + let retry_notify = notify_clone.clone(); + let retry_manager_cancel_token = manager_cancel_token.clone(); + let retry_config = config_for_spawn.clone(); + let retry_panic_id = retry_request_id.clone(); + let retry_panic_state = PanicCleanupState { + active_heals: retry_active_heals.clone(), + heal_queue: retry_heal_queue.clone(), + completed_heals: retry_completed_heals.clone(), + task_aliases: retry_task_aliases.clone(), + retrying_heals: retrying_heals_for_spawn.clone(), + mrf_repair_notice_targets: retry_mrf_repair_notice_targets.clone(), + replacement_recovery_anchors: replacement_recovery_anchors_clone.clone(), + statistics: statistics_clone.clone(), + }; + let retry_panic_cancel_token = retry_cancel_token.clone(); + tokio::spawn(async move { + let retry_child = async move { + #[cfg(test)] + panic_if_armed(SchedulerPanicPoint::RetryChild, &retry_request_id); + loop { + tokio::select! { + _ = retry_cancel_token.cancelled() => { + debug!( + 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() => { + retry_cancel_token.cancel(); + retrying_heals_for_spawn.lock().await.remove(&retry_request_id); + return; + } + _ = sleep(retry_delay) => {} } - retry_completed_heals.lock().await.remove(&retry_request_id); - debug!( - 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" - ); - if should_notify { - retry_notify.notify_one(); + + { + let retrying_heals_guard = retrying_heals_for_spawn.lock().await; + if !retrying_heals_guard.contains_key(&retry_request_id) { + return; + } } - return; - } - HealAdmissionResult::Merged => { - let merged_task_id = - queue.queued_request_id_for_dedup_key(&retry_key).map(ToOwned::to_owned); - retrying_heals_for_spawn.lock().await.remove(&retry_request_id); - drop(queue); - if let Some(merged_task_id) = merged_task_id { + + let active_duplicate_task_id = { + let active_heals_guard = retry_active_heals.lock().await; + active_heal_for_dedup_key(&active_heals_guard, &retry_key).map(|(task_id, _)| task_id) + }; + if let Some(active_duplicate_task_id) = active_duplicate_task_id { + retrying_heals_for_spawn.lock().await.remove(&retry_request_id); move_mrf_repair_notice_targets( &retry_mrf_repair_notice_targets, &retry_request_id, - &merged_task_id, + &active_duplicate_task_id, ); + debug!( + 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; } - debug!( - 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" - ); - return; - } - HealAdmissionResult::Full => { - // admit_request_to_queue already logged the - // rejection (context = "retry"); this repeats - // every backoff cycle while the queue stays - // full, so keep it at debug!. - debug!( - 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_rejected_full", - "Heal retry admission decided" - ); - } - HealAdmissionResult::Dropped(reason) => { - debug!( - 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, - reason = reason.as_str(), - result = "retry_dropped", - "Heal retry admission decided" + + let mut queue = retry_heal_queue.lock().await; + let admission_decision = Self::admit_request_to_queue( + &mut queue, + retry_request.clone(), + &retry_config, + "retry", ); + let admission = admission_decision.result; + let should_notify = matches!(admission, HealAdmissionResult::Accepted) + && retry_config.event_driven_scheduler_enable; + match admission { + HealAdmissionResult::Accepted => { + // Transfer ownership while holding queue -> retrying, + // matching operations_snapshot's lock order. + #[cfg(test)] + pause_retry_ownership_transition(&retry_request_id, true).await; + retrying_heals_for_spawn.lock().await.remove(&retry_request_id); + let displaced_task_id = admission_decision.displaced_task_id; + drop(queue); + if let Some(displaced_task_id) = displaced_task_id { + remove_task_aliases_for_task(&retry_task_aliases, &displaced_task_id).await; + remove_mrf_repair_notice_targets( + &retry_mrf_repair_notice_targets, + &displaced_task_id, + ); + } + retry_completed_heals.lock().await.remove(&retry_request_id); + debug!( + 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" + ); + if should_notify { + retry_notify.notify_one(); + } + return; + } + HealAdmissionResult::Merged => { + let merged_task_id = + queue.queued_request_id_for_dedup_key(&retry_key).map(ToOwned::to_owned); + retrying_heals_for_spawn.lock().await.remove(&retry_request_id); + drop(queue); + if let Some(merged_task_id) = merged_task_id { + move_mrf_repair_notice_targets( + &retry_mrf_repair_notice_targets, + &retry_request_id, + &merged_task_id, + ); + } + debug!( + 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" + ); + return; + } + HealAdmissionResult::Full => { + // admit_request_to_queue already logged the + // rejection (context = "retry"); this repeats + // every backoff cycle while the queue stays + // full, so keep it at debug!. + debug!( + 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_rejected_full", + "Heal retry admission decided" + ); + } + HealAdmissionResult::Dropped(reason) => { + debug!( + 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, + reason = reason.as_str(), + result = "retry_dropped", + "Heal retry admission decided" + ); + } + } } + }; + if AssertUnwindSafe(retry_child).catch_unwind().await.is_err() { + error!( + target: "rustfs::heal::manager", + event = EVENT_HEAL_SCHEDULER_STATE, + component = LOG_COMPONENT_HEAL, + subsystem = LOG_SUBSYSTEM_MANAGER, + request_id = %retry_panic_id, + heal_type = retry_panic_heal_type.kind_label(), + set = %retry_panic_set_label, + state = "retry_child_panicked", + error = PANICKED_HEAL_TASK_ERROR, + "Heal retry child panicked" + ); + finish_panicked_retry_child( + retry_panic_id, + retry_panic_heal_type, + retry_panic_cancel_token, + retry_panic_state, + ) + .await; } - } - }); + }); + } + notify_clone.notify_one(); + }; + if AssertUnwindSafe(scheduler_task).catch_unwind().await.is_err() { + error!( + target: "rustfs::heal::manager", + event = EVENT_HEAL_SCHEDULER_STATE, + component = LOG_COMPONENT_HEAL, + subsystem = LOG_SUBSYSTEM_MANAGER, + task_id = %panic_task_id, + heal_type = panic_task.heal_type.kind_label(), + set = %panic_task.metric_set_label(), + state = "task_panicked", + error = PANICKED_HEAL_TASK_ERROR, + "Heal scheduler task panicked" + ); + finish_panicked_heal_task(panic_task, panic_task_id, panic_state).await; } - notify_clone.notify_one(); }); tasks_started += 1; } else { diff --git a/crates/heal/src/heal/manager/tests.rs b/crates/heal/src/heal/manager/tests.rs index ad50c4bf7..22fc60908 100644 --- a/crates/heal/src/heal/manager/tests.rs +++ b/crates/heal/src/heal/manager/tests.rs @@ -110,7 +110,10 @@ impl HealStorageAPI for MockStorage { Ok(Vec::new()) } - async fn get_bucket_info(&self, _bucket: &str) -> Result> { + async fn get_bucket_info(&self, bucket: &str) -> Result> { + if bucket == "panic" { + panic!("test-only panic payload must not escape the scheduler"); + } Ok(None) } @@ -1021,6 +1024,231 @@ async fn test_task_alias_is_removed_after_terminal_completion() { assert_eq!(manager.canonical_task_id(&duplicate_id).await, duplicate_id); } +#[tokio::test] +#[serial_test::serial] +async fn scheduler_panic_releases_active_slot_and_allows_same_target_readmission() { + let storage: Arc = Arc::new(MockStorage); + let manager = HealManager::new(storage, None); + let request = bucket_request("panic", HealPriority::Normal, HealRequestSource::Admin); + let task_id = request.id.clone(); + + assert_eq!( + manager + .submit_heal_request(request) + .await + .expect("panic request should be admitted"), + HealAdmissionResult::Accepted + ); + let duplicate = bucket_request("panic", HealPriority::Normal, HealRequestSource::Admin); + let duplicate_id = duplicate.id.clone(); + assert_eq!( + manager + .submit_heal_request(duplicate) + .await + .expect("same target should merge while active is queued"), + HealAdmissionResult::Merged + ); + assert_eq!(manager.canonical_task_id(&duplicate_id).await, task_id); + process_manager_queue_once(&manager).await; + + let status = tokio::time::timeout(Duration::from_secs(1), async { + loop { + if let Ok(status) = manager.get_task_status(&task_id).await + && matches!(status, HealTaskStatus::Failed { .. }) + { + break status; + } + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) + .await + .expect("panic task should reach a terminal status"); + assert_eq!( + status, + HealTaskStatus::Failed { + error: PANICKED_HEAL_TASK_ERROR.to_string() + } + ); + assert_eq!(manager.get_active_task_count().await, 0); + assert_eq!(manager.get_queue_length().await, 0); + assert!(manager.retrying_heals.lock().await.is_empty()); + assert!(manager.task_aliases.lock().await.is_empty()); + assert!(manager.completed_heals.lock().await.contains_key(&task_id)); + assert_eq!(manager.canonical_task_id(&duplicate_id).await, duplicate_id); + + let readmitted = bucket_request("panic", HealPriority::Normal, HealRequestSource::Admin); + assert_eq!( + manager + .submit_heal_request(readmitted) + .await + .expect("same target should be re-admitted after a panic"), + HealAdmissionResult::Accepted + ); +} + +#[tokio::test] +#[serial_test::serial] +async fn retry_child_panic_finishes_parent_once() { + clear_scheduler_panic(); + let storage: Arc = Arc::new(MockStorage); + let manager = HealManager::new(storage, None); + let request = HealRequest::object("retry-transition".to_string(), "object".to_string(), None); + let task_id = request.id.clone(); + assert_eq!( + manager + .submit_heal_request(request) + .await + .expect("retry request should be admitted"), + HealAdmissionResult::Accepted + ); + arm_scheduler_panic(SchedulerPanicPoint::RetryChild, &task_id); + process_manager_queue_once(&manager).await; + + let status = tokio::time::timeout(Duration::from_secs(1), async { + loop { + if let Ok(status) = manager.get_task_status(&task_id).await + && matches!(status, HealTaskStatus::Failed { .. }) + { + break status; + } + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) + .await + .expect("retry child panic should finish the parent"); + clear_scheduler_panic(); + assert_eq!( + status, + HealTaskStatus::Failed { + error: PANICKED_HEAL_TASK_ERROR.to_string() + } + ); + assert_eq!(manager.get_active_task_count().await, 0); + assert_eq!(manager.get_queue_length().await, 0); + assert!(manager.retrying_heals.lock().await.is_empty()); + assert!(manager.task_aliases.lock().await.is_empty()); + assert_eq!(manager.completed_heals.lock().await.len(), 1); + assert_eq!(manager.get_statistics().await.failed_tasks, 1); +} + +#[tokio::test] +#[serial_test::serial] +async fn cleanup_panic_is_supervised() { + clear_scheduler_panic(); + let notice_bucket = "cleanup-panic-mrf"; + let notice_object = "object"; + let _ = rustfs_common::mrf_channel::take_mrf_repaired_events_for(notice_bucket); + let storage: Arc = Arc::new(MockStorage); + let manager = HealManager::new(storage, None); + let mut request = HealRequest::new(HealType::Cluster, HealOptions::default(), HealPriority::Normal); + request.source = HealRequestSource::Admin; + let task_id = request.id.clone(); + assert_eq!( + manager + .submit_heal_request(request) + .await + .expect("cleanup request should be admitted"), + HealAdmissionResult::Accepted + ); + manager + .mrf_repair_notice_targets + .lock() + .expect("mrf repair notice registry poisoned") + .insert( + task_id.clone(), + vec![MrfRepairNoticeTarget { + bucket: Arc::from(notice_bucket), + object: Arc::from(notice_object), + version_id: None, + }], + ); + arm_scheduler_panic(SchedulerPanicPoint::Cleanup, &task_id); + process_manager_queue_once(&manager).await; + + let status = tokio::time::timeout(Duration::from_secs(1), async { + loop { + if let Ok(status) = manager.get_task_status(&task_id).await + && matches!(status, HealTaskStatus::Completed) + { + break status; + } + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) + .await + .expect("cleanup panic should leave a terminal status"); + clear_scheduler_panic(); + assert_eq!(status, HealTaskStatus::Completed); + assert_eq!(manager.get_active_task_count().await, 0); + assert!(manager.task_aliases.lock().await.is_empty()); + assert_eq!(manager.completed_heals.lock().await.len(), 1); + assert_eq!(manager.get_statistics().await.successful_tasks, 1); + let events = rustfs_common::mrf_channel::take_mrf_repaired_events_for(notice_bucket); + assert_eq!(events.len(), 1, "cleanup panic must preserve successful MRF notice delivery"); + assert_eq!(events[0].object.as_ref(), notice_object); +} + +#[tokio::test] +#[serial_test::serial] +async fn cancelled_retry_child_panic_does_not_rearchive_failed_status() { + let manager = HealManager::new(Arc::new(MockStorage), None); + let request = HealRequest::object("retry-transition".to_string(), "object".to_string(), None); + let task_id = request.id.clone(); + let retry_cancel_token = insert_retrying_request(&manager, request.clone()).await; + + manager + .cancel_task(&task_id) + .await + .expect("retry cancellation should succeed"); + assert!(retry_cancel_token.is_cancelled()); + + let state = PanicCleanupState { + active_heals: manager.active_heals.clone(), + heal_queue: manager.heal_queue.clone(), + completed_heals: manager.completed_heals.clone(), + task_aliases: manager.task_aliases.clone(), + retrying_heals: manager.retrying_heals.clone(), + mrf_repair_notice_targets: manager.mrf_repair_notice_targets.clone(), + replacement_recovery_anchors: manager.replacement_recovery_anchors.clone(), + statistics: manager.statistics.clone(), + }; + finish_panicked_retry_child(task_id.clone(), request.heal_type, retry_cancel_token, state).await; + + assert!(manager.retrying_heals.lock().await.is_empty()); + assert!(manager.completed_heals.lock().await.is_empty()); + assert_eq!(manager.get_statistics().await.failed_tasks, 0); +} + +#[tokio::test] +async fn active_cancel_wins_parent_panic_cleanup_without_completed_status() { + let manager = HealManager::new(Arc::new(MockStorage), None); + let request = HealRequest::new(HealType::Cluster, HealOptions::default(), HealPriority::Normal); + let task_id = request.id.clone(); + let task = Arc::new(HealTask::from_request(request, Arc::new(MockStorage))); + manager.active_heals.lock().await.insert(task_id.clone(), task.clone()); + + manager + .cancel_task(&task_id) + .await + .expect("active task cancellation should win"); + assert_eq!(task.get_status().await, HealTaskStatus::Cancelled); + + let state = PanicCleanupState { + active_heals: manager.active_heals.clone(), + heal_queue: manager.heal_queue.clone(), + completed_heals: manager.completed_heals.clone(), + task_aliases: manager.task_aliases.clone(), + retrying_heals: manager.retrying_heals.clone(), + mrf_repair_notice_targets: manager.mrf_repair_notice_targets.clone(), + replacement_recovery_anchors: manager.replacement_recovery_anchors.clone(), + statistics: manager.statistics.clone(), + }; + finish_panicked_heal_task(task, task_id, state).await; + + assert!(manager.completed_heals.lock().await.is_empty()); + assert_eq!(manager.get_statistics().await.failed_tasks, 0); +} + #[tokio::test] async fn test_duplicate_admission_is_atomic_with_queue_to_active_transition() { let storage: Arc = Arc::new(MockStorage);