From a9fd0f3c9c0dd19ff6af42e13f89408f2d60d34f 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 18:16:21 +0800 Subject: [PATCH] fix(heal): align scheduler conflict with main --- crates/heal/src/heal/manager/scheduler.rs | 733 ++++------------------ crates/heal/src/heal/manager/tests.rs | 230 +------ 2 files changed, 125 insertions(+), 838 deletions(-) diff --git a/crates/heal/src/heal/manager/scheduler.rs b/crates/heal/src/heal/manager/scheduler.rs index e471444d1..c76a037ae 100644 --- a/crates/heal/src/heal/manager/scheduler.rs +++ b/crates/heal/src/heal/manager/scheduler.rs @@ -13,350 +13,6 @@ // 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>, -} - -async fn cleanup_panicked_task_ownership(state: &PanicCleanupState, task_id: &str) -> bool { - // Queue admission transfers ownership while holding queue -> retrying. Use - // the same order here so panic cleanup cannot remove a retry after another - // request has merged into or admitted that retry. - let mut queue = state.heal_queue.lock().await; - let mut retrying = state.retrying_heals.lock().await; - if retrying.contains_key(task_id) || queue.contains_request_id(task_id) { - publish_heal_queue_length(&queue); - return true; - } - let removed_retry = retrying.remove(task_id).map(|retrying| retrying.cancel_token); - queue.remove_request_id(task_id); - publish_heal_queue_length(&queue); - drop(retrying); - drop(queue); - if let Some(cancel_token) = removed_retry { - cancel_token.cancel(); - } - false -} - -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(cleanup_panicked_task_ownership(&state, &task_id)) - .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 @@ -542,39 +198,54 @@ 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 { - 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!( + 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!( target: "rustfs::heal::manager", event = EVENT_HEAL_SCHEDULER_STATE, component = LOG_COMPONENT_HEAL, @@ -582,8 +253,9 @@ impl HealManager { task_id, heal_type = %task_type_label_for_spawn, set = %task_set_label_for_spawn, - state = "task_completed", - "Heal scheduler task completed" + state = "task_failed", + error = %e, + "Heal scheduler task failed" ); } } @@ -749,74 +421,14 @@ impl HealManager { event = EVENT_HEAL_QUEUE_ADMISSION, 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" + request_id = %retry_request_id, + priority = ?retry_priority, + retry_attempt, + result = "retry_merged_active_duplicate", + "Heal retry admission decided" ); + 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 = @@ -856,183 +468,86 @@ impl HealManager { &displaced_task_id, ); } - - { - let retrying_heals_guard = retrying_heals_for_spawn.lock().await; - if !retrying_heals_guard.contains_key(&retry_request_id) { - return; - } + 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 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); + 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, - &active_duplicate_task_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_active_duplicate", - "Heal retry admission decided" - ); - return; } - - let mut queue = retry_heal_queue.lock().await; - let admission_decision = Self::admit_request_to_queue( - &mut queue, - retry_request.clone(), - &retry_config, - "retry", + 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 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 d81fa3530..585b7a4e7 100644 --- a/crates/heal/src/heal/manager/tests.rs +++ b/crates/heal/src/heal/manager/tests.rs @@ -111,10 +111,7 @@ impl HealStorageAPI for MockStorage { Ok(Vec::new()) } - async fn get_bucket_info(&self, bucket: &str) -> Result> { - if bucket == "panic" { - panic!("test-only panic payload must not escape the scheduler"); - } + async fn get_bucket_info(&self, _bucket: &str) -> Result> { Ok(None) } @@ -1025,231 +1022,6 @@ 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);