Compare commits

..

3 Commits

Author SHA1 Message Date
cxymds e99deef026 Merge branch 'main' into cxymds/fix-1928-scheduler-panic 2026-08-22 11:26:37 +08:00
houseme 2f0918f60b feat(disk): fsync dedicated blocking pool (default-off) (#6366) 2026-08-22 11:24:24 +08:00
马登山 1a6b870eb5 fix(heal): supervise scheduler task panics 2026-08-22 02:39:36 +08:00
5 changed files with 1021 additions and 511 deletions
+7
View File
@@ -57,6 +57,13 @@ pub const DEFAULT_MAX_IO_EVENTS_PER_TICK: usize = 1024;
pub const DEFAULT_EVENT_INTERVAL: u32 = 61;
pub const DEFAULT_RNG_SEED: Option<u64> = None; // None means random
/// Dedicated blocking thread pool for fsync/fdatasync operations.
/// When > 1, fsync operations are isolated from the main blocking pool to
/// prevent device-bound fsync from starving read operations (pread/stat/open).
/// Default 0 means auto (no isolation, use main runtime).
pub const ENV_FSYNC_BLOCKING_THREADS: &str = "RUSTFS_RUNTIME_FSYNC_BLOCKING_THREADS";
pub const DEFAULT_FSYNC_BLOCKING_THREADS: usize = 0;
// Dial9 Tokio Telemetry Default values
pub const DEFAULT_RUNTIME_DIAL9_ENABLED: bool = false; // Disabled by default
pub const DEFAULT_RUNTIME_DIAL9_OUTPUT_DIR: &str = "/var/log/rustfs/telemetry";
+24 -201
View File
@@ -36,8 +36,7 @@ use crate::disk::error::DiskError;
use crate::disk::{BUCKET_META_PREFIX, RUSTFS_META_BUCKET};
use crate::error::{Error, Result};
use crate::error::{
StorageError, is_err_bucket_exists, is_err_bucket_not_found, is_err_object_not_found, is_err_operation_canceled,
is_err_version_not_found,
StorageError, is_err_bucket_exists, is_err_bucket_not_found, is_err_object_not_found, is_err_version_not_found,
};
use crate::layout::endpoints::EndpointServerPools;
use crate::object_api::{GetObjectReader, ObjectOptions};
@@ -774,76 +773,7 @@ async fn load_decommission_entry_exact_versions(
}
fn resolve_decommission_check_after_list_result(list_result: Result<()>, entry_error: Option<Error>) -> Result<()> {
match list_result {
Ok(()) => entry_error.map_or(Ok(()), Err),
Err(list_err) => resolve_decommission_listing_error(Some(list_err), entry_error).map_or(Ok(()), Err),
}
}
fn resolve_decommission_listing_error(listing_error: Option<Error>, entry_error: Option<Error>) -> Option<Error> {
match (listing_error, entry_error) {
(Some(listing_error), Some(entry_error)) if is_err_operation_canceled(&listing_error) => Some(entry_error),
(Some(listing_error), Some(entry_error)) if is_err_operation_canceled(&entry_error) => Some(listing_error),
(Some(listing_error), _) => Some(listing_error),
(None, entry_error) => entry_error,
}
}
fn decommission_unresolved_listing_error(
bucket: &str,
prefix: &str,
candidate: Option<&str>,
candidate_count: usize,
disk_error_count: usize,
pool_index: usize,
set_index: usize,
) -> Error {
let location = candidate.unwrap_or(prefix);
Error::other(format!(
"decommission listing could not resolve metadata for {bucket}/{location} on pool {pool_index} set {set_index} ({candidate_count} candidate(s), {disk_error_count} disk error(s))"
))
}
fn resolve_decommission_partial_listing_entry(
entries: MetaCacheEntries,
resolver: MetadataResolutionParams,
bucket: &str,
prefix: &str,
disk_error_count: usize,
pool_index: usize,
set_index: usize,
) -> Result<MetaCacheEntry> {
let candidate_count = entries.as_ref().iter().flatten().count();
if let Some(entry) = entries.resolve(resolver) {
return Ok(entry);
}
let candidate = entries.as_ref().iter().flatten().map(|entry| entry.name.as_str()).next();
Err(decommission_unresolved_listing_error(
bucket,
prefix,
candidate,
candidate_count,
disk_error_count,
pool_index,
set_index,
))
}
async fn record_decommission_entry_error(
entry_error: &Arc<tokio::sync::Mutex<Option<Error>>>,
rx: &CancellationToken,
err: Error,
) {
if rx.is_cancelled() {
return;
}
let mut first_err = entry_error.lock().await;
if first_err.is_none() && !rx.is_cancelled() {
*first_err = Some(err);
rx.cancel();
}
if let Some(err) = entry_error { Err(err) } else { list_result }
}
fn resolve_decommission_pool_meta_reload_result(result: Result<()>, stage: &str) -> Result<()> {
@@ -3608,7 +3538,6 @@ impl ECStore {
let rx_clone = rx.clone();
let bi = bi.clone();
let set_id = set_idx;
let listing_entry_error = entry_error.clone();
let worker = tokio::spawn(async move {
let _listing_permit = listing_permit;
run_decommission_listing_with_retry(
@@ -3622,11 +3551,7 @@ impl ECStore {
let set = set.clone();
let rx = rx_clone.clone();
let bucket = bi.clone();
let entry_error = listing_entry_error.clone();
async move {
set.list_objects_to_decommission(rx, bucket, callback, entry_error.clone(), idx, set_id)
.await
}
async move { set.list_objects_to_decommission(rx, bucket, callback).await }
},
)
.await
@@ -3656,7 +3581,11 @@ impl ECStore {
wait_decommission_worker_drain(&workers, worker_limit).await?;
if let Some(err) = resolve_decommission_listing_error(listing_worker_error, entry_error.lock().await.clone()) {
if let Some(err) = listing_worker_error {
return Err(err);
}
if let Some(err) = entry_error.lock().await.clone() {
return Err(err);
}
@@ -4262,7 +4191,7 @@ impl ECStore {
let buckets = self.get_buckets_to_decommission().await?;
let pool = self.pools[idx].clone();
for (set_index, set) in pool.disk_set.iter().enumerate() {
for set in &pool.disk_set {
for bucket_info in &buckets {
let mut lifecycle_config = None;
let mut object_lock_config = None;
@@ -4357,7 +4286,7 @@ impl ECStore {
});
let list_result = set
.list_objects_to_decommission(callback_rx, bucket_info.clone(), callback, entry_error.clone(), idx, set_index)
.list_objects_to_decommission(callback_rx, bucket_info.clone(), callback)
.await;
let entry_error = entry_error.lock().await.clone();
resolve_decommission_check_after_list_result(list_result, entry_error)?;
@@ -5092,15 +5021,12 @@ mod tests {
pub type ListCallback = Arc<dyn Fn(MetaCacheEntry) -> BoxFuture<'static, ()> + Send + Sync + 'static>;
impl SetDisks {
#[tracing::instrument(skip(self, rx, cb_func, entry_error))]
#[tracing::instrument(skip(self, rx, cb_func))]
async fn list_objects_to_decommission(
self: &Arc<Self>,
rx: CancellationToken,
bucket_info: DecomBucketInfo,
cb_func: ListCallback,
entry_error: Arc<tokio::sync::Mutex<Option<Error>>>,
pool_index: usize,
set_index: usize,
) -> Result<()> {
let (disks, _) = self.get_online_disks_with_healing(false).await;
ensure_decommission_listing_disks_available(!disks.is_empty(), &bucket_info.name)?;
@@ -5115,12 +5041,6 @@ impl SetDisks {
};
let cb1 = cb_func.clone();
let unresolved_error = entry_error.clone();
let unresolved_rx = rx.clone();
let unresolved_bucket = bucket_info.name.clone();
let unresolved_prefix = bucket_info.prefix.clone();
let unresolved_pool_index = pool_index;
let unresolved_set_index = set_index;
list_path_raw(
rx,
@@ -5133,51 +5053,20 @@ impl SetDisks {
skip_walkdir_total_timeout: true,
walkdir_stall_timeout: Some(DECOMMISSION_BACKGROUND_WALKDIR_STALL_TIMEOUT),
agreed: Some(Box::new(move |entry: MetaCacheEntry| Box::pin(cb1(entry)))),
partial: Some(Box::new(move |entries: MetaCacheEntries, errs: &[Option<DiskError>]| {
partial: Some(Box::new(move |entries: MetaCacheEntries, _: &[Option<DiskError>]| {
let resolver = resolver.clone();
let cb_func = cb_func.clone();
let bucket = unresolved_bucket.clone();
let prefix = unresolved_prefix.clone();
let unresolved_error = unresolved_error.clone();
let unresolved_rx = unresolved_rx.clone();
let pool_index = unresolved_pool_index;
let set_index = unresolved_set_index;
let disk_error_count = errs.iter().flatten().count();
if unresolved_rx.is_cancelled() {
return Box::pin(async {});
}
match resolve_decommission_partial_listing_entry(
entries,
resolver,
&bucket,
&prefix,
disk_error_count,
pool_index,
set_index,
) {
Ok(entry) => {
match entries.resolve(resolver) {
Some(entry) => {
warn!("decommission_pool: list_objects_to_decommission get {}", &entry.name);
Box::pin(async move {
cb_func(entry).await;
})
}
Err(err) => Box::pin(async move {
if unresolved_rx.is_cancelled() {
return;
}
warn!(
event = EVENT_DECOMMISSION_BUCKET,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_POOLS,
bucket = %bucket,
prefix = %prefix,
state = "unresolved_entry",
error = %err,
"Decommission listing failed closed on unresolved metadata"
);
record_decommission_entry_error(&unresolved_error, &unresolved_rx, err).await;
}),
None => {
warn!("decommission_pool: list_objects_to_decommission get none");
Box::pin(async {})
}
}
})),
..Default::default()
@@ -5185,10 +5074,6 @@ impl SetDisks {
)
.await?;
if let Some(err) = entry_error.lock().await.clone() {
return Err(err);
}
Ok(())
}
}
@@ -5394,12 +5279,11 @@ mod pools_tests {
has_active_decommission_canceler, is_decommission_active, is_decommission_cancel_requested,
load_decommission_entry_versions, local_decommission_queue_prefix, mark_decommission_bucket_done,
merge_pool_status_refresh, missing_decommission_worker_prefix, observe_decommission_terminal_reload_result,
pool_meta_has_active_decommission, record_decommission_entry_error, require_decommission_store,
resolve_decommission_bucket_done_save_result, resolve_decommission_bucket_state,
resolve_decommission_check_after_list_result, resolve_decommission_entry_cleanup_delete_result,
resolve_decommission_entry_exact_versions, resolve_decommission_entry_reload_result, resolve_decommission_listing_error,
resolve_decommission_listing_worker_result, resolve_decommission_optional_bucket_config_result,
resolve_decommission_partial_listing_entry, resolve_decommission_pool_meta_reload_result,
pool_meta_has_active_decommission, require_decommission_store, resolve_decommission_bucket_done_save_result,
resolve_decommission_bucket_state, resolve_decommission_check_after_list_result,
resolve_decommission_entry_cleanup_delete_result, resolve_decommission_entry_exact_versions,
resolve_decommission_entry_reload_result, resolve_decommission_listing_worker_result,
resolve_decommission_optional_bucket_config_result, resolve_decommission_pool_meta_reload_result,
resolve_decommission_preflight_heal_result, resolve_decommission_progress_save_result,
resolve_decommission_spawn_failure_result, resolve_decommission_terminal_mark_after_error_result,
resolve_decommission_terminal_mark_result, resolve_decommission_update_after_result,
@@ -5418,9 +5302,7 @@ mod pools_tests {
use crate::error::{Error, StorageError};
use crate::layout::endpoints::{EndpointServerPools, Endpoints, PoolEndpoints};
use crate::services::rebalance::{RebalStatus, RebalanceInfo, RebalanceMeta, RebalanceStats};
use rustfs_filemeta::{
FileInfo, FileInfoVersions, MetaCacheEntries, MetaCacheEntry, MetadataResolutionParams, ObjectPartInfo,
};
use rustfs_filemeta::{FileInfo, FileInfoVersions, MetaCacheEntry, ObjectPartInfo};
use rustfs_rio::Index;
use std::sync::{
Arc,
@@ -6439,65 +6321,6 @@ mod pools_tests {
assert!(matches!(err, Error::SlowDown));
}
#[test]
fn test_resolve_decommission_partial_listing_entry_rejects_unresolved_metadata() {
let err = resolve_decommission_partial_listing_entry(
MetaCacheEntries(vec![None]),
MetadataResolutionParams {
dir_quorum: 2,
obj_quorum: 2,
bucket: "bucket-a".to_string(),
..Default::default()
},
"bucket-a",
"prefix/",
1,
2,
3,
)
.expect_err("unresolved partial listing must fail closed");
let message = err.to_string();
assert!(message.contains("decommission listing could not resolve metadata"));
assert!(message.contains("bucket-a/prefix/"));
assert!(message.contains("pool 2 set 3"));
assert!(message.contains("1 disk error(s)"));
}
#[tokio::test]
async fn test_record_decommission_entry_error_cancels_listing_and_preserves_first_error() {
let entry_error = Arc::new(tokio::sync::Mutex::new(None));
let rx = CancellationToken::new();
record_decommission_entry_error(&entry_error, &rx, Error::SlowDown).await;
record_decommission_entry_error(&entry_error, &rx, Error::OperationCanceled).await;
assert!(rx.is_cancelled());
assert!(matches!(*entry_error.lock().await, Some(Error::SlowDown)));
}
#[tokio::test]
async fn test_record_decommission_entry_error_ignores_already_canceled_listing() {
let entry_error = Arc::new(tokio::sync::Mutex::new(None));
let rx = CancellationToken::new();
rx.cancel();
record_decommission_entry_error(&entry_error, &rx, Error::SlowDown).await;
assert!(entry_error.lock().await.is_none());
}
#[test]
fn test_resolve_decommission_listing_error_preserves_real_listing_failure() {
let err = resolve_decommission_listing_error(Some(Error::SlowDown), Some(Error::OperationCanceled))
.expect("listing failure should be returned");
assert!(matches!(err, Error::SlowDown));
let err = resolve_decommission_listing_error(Some(Error::OperationCanceled), Some(Error::SlowDown))
.expect("entry failure should be returned");
assert!(matches!(err, Error::SlowDown));
}
#[test]
fn test_resolve_decommission_check_after_list_result_returns_list_result_without_entry_error() {
let err = resolve_decommission_check_after_list_result(Err(Error::OperationCanceled), None)
+42 -4
View File
@@ -315,7 +315,7 @@ pub async fn fsync_dir(dir: impl AsRef<Path>) -> io::Result<()> {
#[cfg(unix)]
{
let dir = dir.as_ref().to_path_buf();
tokio::task::spawn_blocking(move || fsync_dir_std(dir)).await?
fsync_spawn_blocking(move || fsync_dir_std(dir)).await?
}
#[cfg(not(unix))]
@@ -683,7 +683,7 @@ async fn fsync_open_dst_dir_group(group: &DstDirFsyncGroup) -> io::Result<()> {
#[cfg(test)]
let dir = group.dir.clone();
let dir_file = group.dir_file.clone();
tokio::task::spawn_blocking(move || {
fsync_spawn_blocking(move || {
#[cfg(test)]
{
if let Some(kind) = fsync_dir_recorder::take_grouped_failure(&dir) {
@@ -1080,6 +1080,44 @@ const TEST_GLOBAL_FILE_SYNCS: usize = 64;
static FILE_SYNC_PERMITS: LazyLock<Semaphore> = LazyLock::new(|| Semaphore::new(global_file_sync_limit()));
static DISK_FILE_SYNC_LIMITERS: LazyLock<Mutex<HashMap<PathBuf, Weak<Semaphore>>>> = LazyLock::new(|| Mutex::new(HashMap::new()));
/// Dedicated tokio runtime for fsync/fdatasync blocking operations. When
/// configured with >1 threads, isolates device-bound fsync from the main
/// blocking pool so reads (pread/stat/open) are not starved. `None` means
/// fall back to the main runtime (zero behavior change).
static FSYNC_RUNTIME: LazyLock<Option<tokio::runtime::Runtime>> = LazyLock::new(|| {
let threads =
rustfs_utils::get_env_usize(rustfs_config::ENV_FSYNC_BLOCKING_THREADS, rustfs_config::DEFAULT_FSYNC_BLOCKING_THREADS);
if threads <= 1 {
return None;
}
let mut builder = tokio::runtime::Builder::new_multi_thread();
builder
.worker_threads(num_cpus::get().min(8))
.max_blocking_threads(threads)
.thread_name("rustfs-fsync")
.thread_stack_size(512 * 1024)
.enable_all();
match builder.build() {
Ok(rt) => {
tracing::info!(threads, "fsync dedicated blocking pool enabled");
Some(rt)
}
Err(err) => {
tracing::warn!(%err, "failed to build fsync runtime, falling back to main pool");
None
}
}
});
/// Spawn a blocking task on the fsync-dedicated runtime if configured,
/// otherwise fall back to the main tokio blocking pool.
fn fsync_spawn_blocking<T: Send + 'static>(f: impl FnOnce() -> T + Send + 'static) -> tokio::task::JoinHandle<T> {
match FSYNC_RUNTIME.as_ref() {
Some(rt) => rt.spawn_blocking(f),
None => tokio::task::spawn_blocking(f),
}
}
static DISK_VOLUME_MUTATION_LOCKS: LazyLock<Mutex<HashMap<PathBuf, Weak<RwLock<()>>>>> =
LazyLock::new(|| Mutex::new(HashMap::new()));
type NamespaceMutationLock = AsyncMutex<()>;
@@ -1217,7 +1255,7 @@ where
F: FnOnce() -> io::Result<T> + Send + 'static,
{
let (disk_permit, global_permit) = acquire_file_sync_permits(disk_permits).await?;
let result = tokio::task::spawn_blocking(move || {
let result = fsync_spawn_blocking(move || {
let _disk_permit = disk_permit;
work()
})
@@ -2146,7 +2184,7 @@ async fn run_blocking_namespace_file_sync_operation_with_global<T: Send + 'stati
wait_started,
);
let disk_permit = admission.disk_permit.clone();
let result = tokio::task::spawn_blocking(move || {
let result = fsync_spawn_blocking(move || {
let _lease = lease;
let _disk_permit = disk_permit;
operation()
File diff suppressed because it is too large Load Diff
+229 -1
View File
@@ -110,7 +110,10 @@ impl HealStorageAPI for MockStorage {
Ok(Vec::new())
}
async fn get_bucket_info(&self, _bucket: &str) -> Result<Option<BucketInfo>> {
async fn get_bucket_info(&self, bucket: &str) -> Result<Option<BucketInfo>> {
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<dyn HealStorageAPI> = 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<dyn HealStorageAPI> = 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<dyn HealStorageAPI> = 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<dyn HealStorageAPI> = Arc::new(MockStorage);