From 2cfefcfd653f7e5338d96334a22e15dfcaf41914 Mon Sep 17 00:00:00 2001 From: cxymds Date: Sun, 13 Sep 2026 10:55:35 +0800 Subject: [PATCH] fix(heal): reject overlapping admin heal owners (#7723) --- crates/config/src/constants/heal.rs | 6 +- crates/heal-contracts/src/heal_channel.rs | 9 +- crates/heal/src/heal/manager.rs | 422 ++++++++------ crates/heal/src/heal/manager/queue.rs | 19 +- crates/heal/src/heal/manager/scheduler.rs | 15 +- crates/heal/src/heal/manager/tests.rs | 8 +- .../src/heal/manager/tests/admin_overlap.rs | 551 ++++++++++++++++++ .../src/heal/manager/tests/root_recovery.rs | 159 ++++- crates/protos/src/heal_control.rs | 6 +- docs/operations/scanner-runtime-controls.md | 2 +- rustfs/src/admin/handlers/heal.rs | 1 + rustfs/src/storage/rpc/node_service.rs | 81 +++ 12 files changed, 1082 insertions(+), 197 deletions(-) create mode 100644 crates/heal/src/heal/manager/tests/admin_overlap.rs diff --git a/crates/config/src/constants/heal.rs b/crates/config/src/constants/heal.rs index 9c27a663f..22bdbdde9 100644 --- a/crates/config/src/constants/heal.rs +++ b/crates/config/src/constants/heal.rs @@ -208,9 +208,9 @@ pub const DEFAULT_HEAL_MRF_REPLAY_BATCH: usize = 256; /// Environment variable selecting how admin heal starts behave when the /// requested path overlaps an already running or queued heal: `merge` -/// (default, keep today's dedup/merge semantics) or `minio_error` (return a -/// typed already-running / overlapping-paths rejection like madmin). +/// (default, merge equivalent starts and reject other overlapping admin +/// scopes) or `minio_error` (also reject equivalent starts like madmin). pub const ENV_HEAL_OVERLAP_POLICY: &str = "RUSTFS_HEAL_OVERLAP_POLICY"; -/// Default overlap policy: merge duplicate/overlapping requests. +/// Default overlap policy: merge equivalent requests, reject overlapping scopes. pub const DEFAULT_HEAL_OVERLAP_POLICY: &str = "merge"; diff --git a/crates/heal-contracts/src/heal_channel.rs b/crates/heal-contracts/src/heal_channel.rs index 2e4f9a056..cef8fa0ab 100644 --- a/crates/heal-contracts/src/heal_channel.rs +++ b/crates/heal-contracts/src/heal_channel.rs @@ -225,12 +225,11 @@ pub struct HealOpts { pub enum HealAdmissionDropReason { QueueFull, PolicyDropped, - /// HS-06: an admin heal start overlaps (same bucket with mutually - /// containing prefixes, or the same erasure set) an already running or - /// queued task. Only produced when RUSTFS_HEAL_OVERLAP_POLICY=minio_error. + /// An admin target already has an incompatible or durable-only owner, + /// or an equivalent start was rejected by the `minio_error` policy. AlreadyRunning, - /// HS-06: same as [`Self::AlreadyRunning`] but for paths that merely - /// contain (or are contained by) the active task's path. + /// An admin scope intersects another owner's scope without being an + /// equivalent request. Produced by both overlap policies. OverlappingPaths, } diff --git a/crates/heal/src/heal/manager.rs b/crates/heal/src/heal/manager.rs index 4540aec88..c76a0f74e 100644 --- a/crates/heal/src/heal/manager.rs +++ b/crates/heal/src/heal/manager.rs @@ -515,7 +515,7 @@ fn publish_heal_queue_length(queue: &PriorityHealQueue) { fn active_heal_for_dedup_key(active_heals: &HashMap>, key: &str) -> Option<(String, HealType)> { active_heals .iter() - .find(|(_, task)| PriorityHealQueue::make_dedup_key_for_type(&task.heal_type) == key) + .find(|(_, task)| PriorityHealQueue::make_dedup_key_for_scope(&task.heal_type, &task.options) == key) .map(|(task_id, task)| (task_id.clone(), task.heal_type.clone())) } @@ -608,8 +608,7 @@ fn recoverable_heal_retry_delay(retry_attempt: u32) -> Duration { /// HS-06 admin overlap policy. #[derive(Debug, Clone, Copy, PartialEq, Eq, Default)] pub enum HealOverlapPolicy { - /// Default: overlapping admin starts merge into the existing task - /// (today's dedup semantics). + /// Merge equivalent admin starts; reject different overlapping admin scopes. #[default] Merge, /// Return a typed already-running / overlapping-paths rejection like @@ -617,19 +616,16 @@ pub enum HealOverlapPolicy { MinioError, } -/// Path view of a heal type for overlap comparison: a bucket plus a -/// prefix/object path inside it (`None` bucket = cluster-wide, overlaps -/// everything). -fn heal_type_path_view(heal_type: &HealType) -> (Option<&str>, &str) { +/// S3 key bytes are already decoded at the request boundary. An object name +/// denotes one key, while a prefix (including the bucket's empty prefix) is a range. +fn heal_type_path_view(heal_type: &HealType) -> Option<(&str, &str, bool)> { match heal_type { - HealType::Cluster => (None, ""), - HealType::Bucket { bucket } => (Some(bucket), ""), - HealType::Prefix { bucket, prefix } => (Some(bucket), prefix), + HealType::Bucket { bucket } => Some((bucket, "", true)), + HealType::Prefix { bucket, prefix } => Some((bucket, prefix, true)), HealType::Object { bucket, object, .. } | HealType::Metadata { bucket, object } - | HealType::ECDecode { bucket, object, .. } => (Some(bucket), object), - // Erasure-set heal: the set id is the overlap dimension. - HealType::ErasureSet { set_disk_id, .. } => (Some("\u{0}set"), set_disk_id), + | HealType::ECDecode { bucket, object, .. } => Some((bucket, object, false)), + HealType::Cluster | HealType::ErasureSet { .. } => None, } } @@ -644,35 +640,88 @@ enum OverlapVerdict { Overlapping, } -fn prefix_paths_overlap(a: &str, b: &str) -> OverlapVerdict { - if a == b { - return OverlapVerdict::SameTarget; +fn heal_scope_indices(heal_type: &HealType, options: &HealOptions) -> (Option, Option) { + if let HealType::ErasureSet { set_disk_id, .. } = heal_type + && let Ok((pool, set)) = super::utils::parse_set_disk_id(set_disk_id) + { + return (Some(pool), Some(set)); } - if a.is_empty() || b.is_empty() || a.starts_with(b) || b.starts_with(a) { - return OverlapVerdict::Overlapping; - } - OverlapVerdict::Disjoint + (options.pool_index, options.set_index) } -fn heal_types_overlap(left: &HealType, right: &HealType) -> OverlapVerdict { - let (left_bucket, left_path) = heal_type_path_view(left); - let (right_bucket, right_path) = heal_type_path_view(right); - match (left_bucket, right_bucket) { - // Cluster-wide overlaps everything (but an exact cluster match is - // SameTarget). - (None, _) | (_, None) => { - if matches!(left, HealType::Cluster) && matches!(right, HealType::Cluster) { - OverlapVerdict::SameTarget - } else { - OverlapVerdict::Overlapping - } +fn heal_scopes_overlap( + left: &HealType, + left_options: &HealOptions, + right: &HealType, + right_options: &HealOptions, +) -> OverlapVerdict { + let left_scope = heal_scope_indices(left, left_options); + let right_scope = heal_scope_indices(right, right_options); + if matches!((left_scope.0, right_scope.0), (Some(a), Some(b)) if a != b) + || matches!((left_scope.1, right_scope.1), (Some(a), Some(b)) if a != b) + { + return OverlapVerdict::Disjoint; + } + + let overlaps = match (left, right) { + (HealType::Cluster, _) | (_, HealType::Cluster) => true, + // Erasure-set tasks also repair the set's structure, even if their + // bucket lists are disjoint. An empty list means every bucket. + (HealType::ErasureSet { .. }, HealType::ErasureSet { .. }) => true, + (HealType::ErasureSet { buckets, .. }, other) | (other, HealType::ErasureSet { buckets, .. }) => { + buckets.is_empty() + || heal_type_path_view(other).is_some_and(|(bucket, _, _)| buckets.iter().any(|candidate| candidate == bucket)) } - (Some(lb), Some(rb)) => { - if lb != rb { - return OverlapVerdict::Disjoint; - } - prefix_paths_overlap(left_path, right_path) + _ => { + matches!((heal_type_path_view(left), heal_type_path_view(right)), + (Some((lb, lp, left_prefix)), Some((rb, rp, right_prefix))) + if lb == rb && (lp == rp || (left_prefix && rp.starts_with(lp)) || (right_prefix && lp.starts_with(rp)))) } + }; + if !overlaps { + OverlapVerdict::Disjoint + } else if left == right && left_scope == right_scope { + OverlapVerdict::SameTarget + } else { + // Versions of the same object share xl.meta and are not independent + // owners, but different versions must never merge into one receipt. + OverlapVerdict::Overlapping + } +} + +fn admin_heal_options_compatible(left: &HealOptions, right: &HealOptions) -> bool { + // The scheduler fills and consumes the timeout budget. It is execution + // state, not a reason to discard an otherwise equivalent owner's token. + left.scan_mode == right.scan_mode + && left.remove_corrupted == right.remove_corrupted + && left.recreate_missing == right.recreate_missing + && left.update_parity == right.update_parity + && left.recursive == right.recursive + && left.dry_run == right.dry_run + && left.no_lock == right.no_lock +} + +fn admin_overlap_rejection( + request: &HealRequest, + owner_type: &HealType, + owner_options: &HealOptions, + owner_source: HealRequestSource, + policy: HealOverlapPolicy, +) -> Option { + if request.source != HealRequestSource::Admin + || (policy == HealOverlapPolicy::Merge && owner_source != HealRequestSource::Admin) + { + return None; + } + match heal_scopes_overlap(&request.heal_type, &request.options, owner_type, owner_options) { + OverlapVerdict::Disjoint => None, + OverlapVerdict::SameTarget + if policy == HealOverlapPolicy::Merge && admin_heal_options_compatible(&request.options, owner_options) => + { + None + } + OverlapVerdict::SameTarget => Some(HealAdmissionDropReason::AlreadyRunning), + OverlapVerdict::Overlapping => Some(HealAdmissionDropReason::OverlappingPaths), } } @@ -696,8 +745,8 @@ pub struct HealConfig { pub low_priority_drop_when_full: bool, /// Whether notify-driven scheduler wakeups are enabled. pub event_driven_scheduler_enable: bool, - /// How admin heal starts behave on path overlap (HS-06): merge into the - /// existing task (default) or return a typed already-running rejection. + /// Merge equivalent admin starts by default; reject different overlapping + /// scopes. The strict policy also rejects equivalent requests. pub overlap_policy: HealOverlapPolicy, /// Whether per-set bulkhead scheduling is enabled. pub set_bulkhead_enable: bool, @@ -857,8 +906,8 @@ pub struct HealManager { replacement_recovery_blocked_sets: Arc>>, /// Durable handoff of interrupted administrator root traversals. root_recovery: Arc, - /// Keep forceStart's cancellation side effects inside the shutdown fence. - force_start_shutdown: Mutex<()>, + /// Serialize admin starts and forceStart replacement with shutdown. + admin_start_shutdown: Mutex<()>, /// Storage layer interface storage: Arc, /// Cancel token @@ -1459,7 +1508,7 @@ impl HealManager { replacement_recovery_anchors: Arc::new(std::sync::Mutex::new(HashMap::new())), replacement_recovery_blocked_sets: Arc::new(std::sync::Mutex::new(HashSet::new())), root_recovery, - force_start_shutdown: Mutex::new(()), + admin_start_shutdown: Mutex::new(()), storage, cancel_token: CancellationToken::new(), statistics: Arc::new(RwLock::new(HealStatistics::new())), @@ -1571,7 +1620,7 @@ impl HealManager { /// Stop HealManager pub async fn stop(&self) -> Result<()> { - let _force_start_guard = self.force_start_shutdown.lock().await; + let _admin_start_guard = self.admin_start_shutdown.lock().await; info!( target: "rustfs::heal::manager", event = EVENT_HEAL_MANAGER_STATE, @@ -1740,10 +1789,10 @@ impl HealManager { let admission_start = Instant::now(); let source = request.source; let force_start = request.force_start; - // A forceStart must not retire an old durable owner if shutdown will - // reject its replacement. Hold the same gate through final admission. - let _force_start_guard = if source == HealRequestSource::Admin && force_start { - let guard = self.force_start_shutdown.lock().await; + // Keep ordinary STARTs outside forceStart's cancel-then-admit window. + // Shutdown uses the same gate, before the active -> queue -> retry locks. + let _admin_start_guard = if source == HealRequestSource::Admin { + let guard = self.admin_start_shutdown.lock().await; if self.cancel_token.is_cancelled() { return Err(Error::Other("Heal manager is stopping".to_string())); } @@ -1751,62 +1800,13 @@ impl HealManager { } else { None }; - // HS-06 forceStart semantics (admin only): MinIO stops the old task - // first and then starts the new one. Cancel any active admin task - // overlapping this request's path before entering admission, so the - // fresh task is never merged into the one being replaced. - if request.source == HealRequestSource::Admin && request.force_start { - let overlapping: Vec = { - let active_heals = self.active_heals.lock().await; - let queue = self.heal_queue.lock().await; - let retrying = self.retrying_heals.lock().await; - let mut ids = active_heals - .iter() - .filter(|(task_id, task)| { - task.source == HealRequestSource::Admin - && heal_types_overlap(&request.heal_type, &task.heal_type) != OverlapVerdict::Disjoint - && *task_id != &request.id - }) - .map(|(task_id, _)| task_id.clone()) - .collect::>(); - ids.extend( - queue - .requests() - .chain(retrying.values().map(|retrying| &retrying.request)) - .filter(|pending| { - root_recovery::is_admin_heal_recovery(&pending.heal_type, pending.source) - && heal_types_overlap(&request.heal_type, &pending.heal_type) != OverlapVerdict::Disjoint - && pending.id != request.id - }) - .map(|pending| pending.id.clone()), - ); - ids - }; - for task_id in overlapping { - match self.cancel_task(&task_id).await { - Ok(_) => info!( - target: "rustfs::heal::manager", - event = EVENT_HEAL_QUEUE_ADMISSION, - component = LOG_COMPONENT_HEAL, - subsystem = LOG_SUBSYSTEM_MANAGER, - request_id = %request.id, - cancelled_task_id = %task_id, - result = "force_start_cancelled_overlap", - "Admin forceStart cancelled an overlapping heal task" - ), - Err(err) => return Err(err), - } - } - // A failed or timed-out replay may have only its durable owner - // left. Cancel only records that overlap this forced start. - for pending in self.root_recovery.pending().await? { - if pending.id != request.id - && heal_types_overlap(&request.heal_type, &pending.heal_type) != OverlapVerdict::Disjoint - { - self.cancel_task(&pending.id).await?; - } - } - } + // Decode all durable responsibilities before forceStart has side effects. + // Do not hold runtime state locks while scanning recovery records. + let mut durable_owners = if source == HealRequestSource::Admin { + self.root_recovery.pending().await? + } else { + Vec::new() + }; let config = self.config.read().await; let dedup_key = PriorityHealQueue::make_dedup_key(&request); @@ -1815,14 +1815,14 @@ impl HealManager { // in the same atomic view. Otherwise queue -> active and // active -> retrying transitions can slip between duplicate checks. let lock_phase_start = Instant::now(); - let active_heals = self.active_heals.lock().await; + let mut active_heals = self.active_heals.lock().await; if self.cancel_token.is_cancelled() { return Err(Error::Other("Heal manager is stopping".to_string())); } #[cfg(test)] pause_duplicate_admission_after_active_lock(&request.id).await; let mut queue = self.heal_queue.lock().await; - let retrying_heals = self.retrying_heals.lock().await; + let mut retrying_heals = self.retrying_heals.lock().await; let request_id_admission = active_heals .get(&request.id) @@ -1837,6 +1837,14 @@ impl HealManager { retrying_heals .get(&request.id) .map(|retrying| (request_matches_request(&request, &retrying.request), "retrying")) + }) + // A durable-only owner must be resumed by recovery, not overwritten + // with a new execution budget by a replay before recovery completes. + .or_else(|| { + durable_owners + .iter() + .find(|owner| owner.id == request.id) + .map(|_| (false, "durable")) }); if let Some((matches_existing, duplicate_state)) = request_id_admission { let admission = if matches_existing { @@ -1885,7 +1893,96 @@ impl HealManager { }); } - let duplicate = (!request.force_start).then(|| { + if source == HealRequestSource::Admin && force_start { + let overlapping = active_heals + .values() + .map(|task| (&task.id, &task.heal_type, &task.options, task.source)) + .chain( + queue + .requests() + .chain(retrying_heals.values().map(|retrying| &retrying.request)) + .chain(durable_owners.iter()) + .map(|owner| (&owner.id, &owner.heal_type, &owner.options, owner.source)), + ) + .filter(|(id, heal_type, options, source)| { + **id != request.id + && *source == HealRequestSource::Admin + && heal_scopes_overlap(&request.heal_type, &request.options, heal_type, options) + != OverlapVerdict::Disjoint + }) + .map(|(id, _, _, _)| id.clone()) + .collect::>(); + drop(retrying_heals); + drop(queue); + drop(active_heals); + for task_id in &overlapping { + match self.cancel_task(task_id).await { + Ok(_) => info!( + target: "rustfs::heal::manager", + event = EVENT_HEAL_QUEUE_ADMISSION, + component = LOG_COMPONENT_HEAL, + subsystem = LOG_SUBSYSTEM_MANAGER, + request_id = %request.id, + cancelled_task_id = %task_id, + result = "force_start_cancelled_overlap", + "Admin forceStart cancelled an overlapping heal task" + ), + Err(Error::TaskNotFound { .. }) => {} // The owner completed before cancellation acquired its state. + Err(err) => return Err(err), + } + } + durable_owners.retain(|owner| !overlapping.contains(&owner.id)); + active_heals = self.active_heals.lock().await; + queue = self.heal_queue.lock().await; + retrying_heals = self.retrying_heals.lock().await; + } + + let rejection = if source == HealRequestSource::Admin { + let live_owners = active_heals + .values() + .map(|task| (task.id.as_str(), &task.heal_type, &task.options, task.source)) + .chain( + queue + .requests() + .chain(retrying_heals.values().map(|retrying| &retrying.request)) + .map(|owner| (owner.id.as_str(), &owner.heal_type, &owner.options, owner.source)), + ) + .collect::>(); + let live_ids = live_owners.iter().map(|(id, _, _, _)| *id).collect::>(); + live_owners + .iter() + .filter_map(|(id, heal_type, options, owner_source)| { + admin_overlap_rejection(&request, heal_type, options, *owner_source, config.overlap_policy) + .map(|reason| (reason, *id)) + }) + .chain( + durable_owners + .iter() + .filter(|owner| !live_ids.contains(owner.id.as_str())) + .filter_map(|owner| { + // A persisted-only task still owns its scope, but cannot serve + // a merged runtime receipt until recovery has restored it. + admin_overlap_rejection( + &request, + &owner.heal_type, + &owner.options, + owner.source, + HealOverlapPolicy::MinioError, + ) + .map(|reason| (reason, owner.id.as_str())) + }), + ) + .min_by(|(left_reason, left_id), (right_reason, right_id)| { + (*left_reason != HealAdmissionDropReason::AlreadyRunning) + .cmp(&(*right_reason != HealAdmissionDropReason::AlreadyRunning)) + .then_with(|| left_id.cmp(right_id)) + }) + .map(|(reason, id)| (reason, id.to_string())) + } else { + None + }; + + let duplicate = (!request.force_start && rejection.is_none()).then(|| { active_heal_for_dedup_key(&active_heals, &dedup_key) .map(|(task_id, _)| (task_id, "active")) .or_else(|| { @@ -1966,70 +2063,36 @@ impl HealManager { }); } - // HS-06 typed overlap rejection (admin only, minio_error policy): - // paths containing or contained by an active/queued task reject with - // AlreadyRunning / OverlappingPaths instead of merging. Exact - // duplicates already merged above; scanner/autoheal/read-repair - // sources never take this path. - if request.source == HealRequestSource::Admin && config.overlap_policy == HealOverlapPolicy::MinioError { - let mut rejection = None; - for (task_id, task) in active_heals.iter() { - match heal_types_overlap(&request.heal_type, &task.heal_type) { - OverlapVerdict::SameTarget => { - rejection = Some((HealAdmissionDropReason::AlreadyRunning, task_id.clone())); - break; - } - OverlapVerdict::Overlapping => { - rejection = Some((HealAdmissionDropReason::OverlappingPaths, task_id.clone())); - } - OverlapVerdict::Disjoint => {} - } - } - if rejection.is_none() { - for queued in queue.requests() { - match heal_types_overlap(&request.heal_type, &queued.heal_type) { - OverlapVerdict::SameTarget => { - rejection = Some((HealAdmissionDropReason::AlreadyRunning, queued.id.clone())); - break; - } - OverlapVerdict::Overlapping => { - rejection = Some((HealAdmissionDropReason::OverlappingPaths, queued.id.clone())); - } - OverlapVerdict::Disjoint => {} - } - } - } - if let Some((reason, overlap_task_id)) = rejection { - drop(retrying_heals); - drop(queue); - drop(active_heals); - let lock_phase = lock_phase_start.elapsed(); - Self::record_admission_metric(request.source, HealAdmissionResult::Dropped(reason), "overlap_rejected"); - self.record_admission_observation(HealAdmissionObservation { - source, - result: HealAdmissionResult::Dropped(reason), - context: "overlap_rejected", - force_start, - displaced: false, - start_duration: admission_start.elapsed(), - lock_phase, - }); - warn!( - target: "rustfs::heal::manager", - event = EVENT_HEAL_QUEUE_ADMISSION, - component = LOG_COMPONENT_HEAL, - subsystem = LOG_SUBSYSTEM_MANAGER, - request_id = %request.id, - overlap_task_id = %overlap_task_id, - reason = reason.as_str(), - result = "overlap_rejected", - "Admin heal start rejected by overlap policy" - ); - return Ok(HealAdmissionReceipt { - result: HealAdmissionResult::Dropped(reason), - task_id: overlap_task_id, - }); - } + if let Some((reason, overlap_task_id)) = rejection { + drop(retrying_heals); + drop(queue); + drop(active_heals); + let lock_phase = lock_phase_start.elapsed(); + Self::record_admission_metric(request.source, HealAdmissionResult::Dropped(reason), "overlap_rejected"); + self.record_admission_observation(HealAdmissionObservation { + source, + result: HealAdmissionResult::Dropped(reason), + context: "overlap_rejected", + force_start, + displaced: false, + start_duration: admission_start.elapsed(), + lock_phase, + }); + warn!( + target: "rustfs::heal::manager", + event = EVENT_HEAL_QUEUE_ADMISSION, + component = LOG_COMPONENT_HEAL, + subsystem = LOG_SUBSYSTEM_MANAGER, + request_id = %request.id, + overlap_task_id = %overlap_task_id, + reason = reason.as_str(), + result = "overlap_rejected", + "Admin heal start rejected by overlap policy" + ); + return Ok(HealAdmissionReceipt { + result: HealAdmissionResult::Dropped(reason), + task_id: overlap_task_id, + }); } let durable_handoff = root_recovery::is_admin_heal_recovery(&request.heal_type, request.source); @@ -2385,8 +2448,12 @@ impl HealManager { /// Cancel task pub async fn cancel_task(&self, task_id: &str) -> Result<()> { let canonical_task_id = self.canonical_task_id(task_id).await; + // Select and retire the owner atomically with scheduler transitions. + // All multi-map paths acquire active -> queue -> retrying. + let mut active_heals = self.active_heals.lock().await; + let mut queue = self.heal_queue.lock().await; + let mut retrying_heals = self.retrying_heals.lock().await; { - let mut active_heals = self.active_heals.lock().await; if let Some(task) = active_heals.get(&canonical_task_id) { let completed = CompletedHealStatus::snapshot(task, HealTaskStatus::Cancelled).await; self.publish_admin_terminal(&canonical_task_id, &task.heal_type, task.source, &completed) @@ -2405,6 +2472,8 @@ impl HealManager { state = "cancelled_active_task", "Heal manager cancelled active task" ); + drop(retrying_heals); + drop(queue); drop(active_heals); self.remove_aliases_for_task(&canonical_task_id).await; self.remove_mrf_repair_notice_targets_for_task(&canonical_task_id); @@ -2413,7 +2482,6 @@ impl HealManager { } { - let mut retrying_heals = self.retrying_heals.lock().await; if let Some(retrying) = retrying_heals.get(&canonical_task_id) { self.publish_admin_cancelled_terminal( &canonical_task_id, @@ -2429,6 +2497,8 @@ impl HealManager { if let Some(retrying) = retrying_heals.remove(&canonical_task_id) { retrying.cancel_token.cancel(); drop(retrying_heals); + drop(queue); + drop(active_heals); self.completed_heals.lock().await.remove(&canonical_task_id); self.remove_aliases_for_task(&canonical_task_id).await; self.remove_mrf_repair_notice_targets_for_task(&canonical_task_id); @@ -2445,7 +2515,6 @@ impl HealManager { } } - let mut queue = self.heal_queue.lock().await; if let Some(request) = queue.requests().find(|request| request.id == canonical_task_id) { self.publish_admin_cancelled_terminal(&canonical_task_id, &request.heal_type, request.source, &request.options) .await?; @@ -2464,14 +2533,19 @@ impl HealManager { state = "cancelled_queued_task", "Heal manager cancelled queued task" ); + drop(retrying_heals); drop(queue); + drop(active_heals); self.remove_aliases_for_task(&canonical_task_id).await; self.remove_mrf_repair_notice_targets_for_task(&canonical_task_id); return Ok(()); } + let cancelled_pending = self.root_recovery.cancel_pending(&canonical_task_id).await?; + drop(retrying_heals); drop(queue); - if self.root_recovery.cancel_pending(&canonical_task_id).await? { + drop(active_heals); + if cancelled_pending { return Ok(()); } Err(Error::TaskNotFound { diff --git a/crates/heal/src/heal/manager/queue.rs b/crates/heal/src/heal/manager/queue.rs index 2fa9962ee..c5b01c55d 100644 --- a/crates/heal/src/heal/manager/queue.rs +++ b/crates/heal/src/heal/manager/queue.rs @@ -435,10 +435,21 @@ impl PriorityHealQueue { /// Create a deduplication key from a heal request pub(super) fn make_dedup_key(request: &HealRequest) -> String { - let base = Self::make_dedup_key_for_type(&request.heal_type); - match (&request.heal_type, request.options.set_key()) { - (HealType::Object { .. } | HealType::ECDecode { .. }, Some(scope)) => format!("{base}:scope:{scope}"), - _ => base, + Self::make_dedup_key_for_scope(&request.heal_type, &request.options) + } + + pub(super) fn make_dedup_key_for_scope(heal_type: &HealType, options: &HealOptions) -> String { + let base = Self::make_dedup_key_for_type(heal_type); + // Erasure-set keys already encode pool/set and are also queried by + // automatic replacement admission through contains_erasure_set. + if matches!(heal_type, HealType::ErasureSet { .. }) { + return base; + } + match heal_scope_indices(heal_type, options) { + (None, None) => base, + // A distinct leading tag cannot alias an unscoped S3 key that + // happens to contain the scope suffix as literal object bytes. + (pool, set) => format!("scope:{pool:?}:{set:?}:{base}"), } } diff --git a/crates/heal/src/heal/manager/scheduler.rs b/crates/heal/src/heal/manager/scheduler.rs index dfe96dfaa..6a05e913a 100644 --- a/crates/heal/src/heal/manager/scheduler.rs +++ b/crates/heal/src/heal/manager/scheduler.rs @@ -505,7 +505,16 @@ impl HealManager { return; } + #[cfg(test)] + tests::admin_overlap::pause_before_retry_queue(&retry_request_id).await; let mut queue = retry_heal_queue.lock().await; + let mut retrying = retrying_heals_for_spawn.lock().await; + // Cancellation may win after the backoff checks but + // before queue acquisition. Keep ownership through + // publication so a cancelled retry cannot reappear. + if retry_cancel_token.is_cancelled() || !retrying.contains_key(&retry_request_id) { + return; + } let admission_decision = Self::admit_request_to_queue(&mut queue, retry_request.clone(), &retry_config, "retry"); let admission = admission_decision.result; @@ -525,7 +534,8 @@ impl HealManager { // 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); + retrying.remove(&retry_request_id); + drop(retrying); let displaced_task_id = admission_decision.displaced_task_id().map(ToOwned::to_owned); drop(queue); if let (Some(displaced_task_id), Some(displaced_terminal)) = @@ -565,7 +575,8 @@ impl HealManager { 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); + retrying.remove(&retry_request_id); + drop(retrying); drop(queue); if let Some(merged_task_id) = merged_task_id { move_mrf_repair_notice_targets( diff --git a/crates/heal/src/heal/manager/tests.rs b/crates/heal/src/heal/manager/tests.rs index d340daded..20b795160 100644 --- a/crates/heal/src/heal/manager/tests.rs +++ b/crates/heal/src/heal/manager/tests.rs @@ -26,6 +26,7 @@ use rustfs_madmin::heal_commands::HealResultItem; use std::sync::Mutex as StdMutex; use tempfile::TempDir; +pub(super) mod admin_overlap; mod root_recovery; mod running_mainline; @@ -3073,17 +3074,16 @@ async fn overlap_policy_minio_error_rejects_same_and_containing_paths() { } #[tokio::test] -async fn overlap_policy_default_merge_keeps_today_semantics() { +async fn overlap_policy_default_merge_rejects_nested_admin_but_preserves_scanner_admission() { let manager = manager_with_policy(HealOverlapPolicy::Merge); insert_active_task(&manager, admin_prefix_request("bucket-a", "logs/")).await; - // Different-dedup-key overlap still merges under the default policy: - // the nested path dedups to its own key but nothing rejects it. + // A different key must not create a second owner of an admin range. let nested = manager .submit_heal_request(admin_prefix_request("bucket-a", "logs/app/")) .await .expect("admission must decide"); - assert_eq!(nested, HealAdmissionResult::Accepted, "default policy must not reject overlaps"); + assert_eq!(nested, HealAdmissionResult::Dropped(HealAdmissionDropReason::OverlappingPaths)); // Non-admin sources never get overlap rejections even under minio_error. let manager = manager_with_policy(HealOverlapPolicy::MinioError); diff --git a/crates/heal/src/heal/manager/tests/admin_overlap.rs b/crates/heal/src/heal/manager/tests/admin_overlap.rs new file mode 100644 index 000000000..703e74f0a --- /dev/null +++ b/crates/heal/src/heal/manager/tests/admin_overlap.rs @@ -0,0 +1,551 @@ +// Copyright 2026 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use super::*; +use rustfs_heal_contracts::heal_channel::HealScanMode; +use std::future::Future; +use std::task::Poll; + +#[derive(Clone, Copy, Debug)] +enum OwnerState { + Active, + Queued, + Retrying, +} + +async fn install_owner(manager: &HealManager, request: HealRequest, state: OwnerState) -> String { + let id = request.id.clone(); + match state { + OwnerState::Active => { + insert_active_task(manager, request).await; + } + OwnerState::Queued => { + assert_eq!(manager.heal_queue.lock().await.push(request), QueuePushOutcome::Accepted); + } + OwnerState::Retrying => { + manager.retrying_heals.lock().await.insert( + id.clone(), + RetryingHeal { + request, + error: "recoverable fixture error".to_string(), + cancel_token: CancellationToken::new(), + }, + ); + } + } + id +} + +fn admin_object_request(object: &str) -> HealRequest { + let mut request = HealRequest::object("bucket".to_string(), object.to_string(), None); + request.source = HealRequestSource::Admin; + request +} + +#[tokio::test] +async fn admin_overlap_default_rejects_parent_child_without_new_owner() { + for (existing, incoming) in [("scope/", "scope/child/"), ("scope/child/", "scope/")] { + let manager = manager_with_policy(HealOverlapPolicy::Merge); + let owner = insert_active_task(&manager, admin_prefix_request("bucket", existing)).await; + let receipt = manager + .submit_heal_request_with_receipt(admin_prefix_request("bucket", incoming)) + .await + .expect("overlapping admin admission should return a typed decision"); + assert_eq!( + receipt.result, + HealAdmissionResult::Dropped(HealAdmissionDropReason::OverlappingPaths), + "existing={existing}, incoming={incoming}" + ); + assert_eq!(receipt.task_id, owner, "a rejection must identify the existing owner"); + assert!(manager.heal_queue.lock().await.is_empty(), "rejected starts must not enter the queue"); + assert_eq!(manager.active_heals.lock().await.len(), 1); + } +} + +#[tokio::test] +async fn admin_overlap_policy_matrix_covers_every_live_owner_state() { + for state in [OwnerState::Active, OwnerState::Queued, OwnerState::Retrying] { + for policy in [HealOverlapPolicy::Merge, HealOverlapPolicy::MinioError] { + for (existing, incoming, relation) in [ + ("scope/", "scope/", OverlapVerdict::SameTarget), + ("scope/", "scope/child/", OverlapVerdict::Overlapping), + ("scope/child/", "scope/", OverlapVerdict::Overlapping), + ("scope/child/", "scope/other/", OverlapVerdict::Disjoint), + ] { + let manager = manager_with_policy(policy); + let owner = install_owner(&manager, admin_prefix_request("bucket", existing), state).await; + let request = admin_prefix_request("bucket", incoming); + let request_id = request.id.clone(); + let receipt = manager + .submit_heal_request_with_receipt(request) + .await + .expect("typed admission"); + let expected = match (relation, policy) { + (OverlapVerdict::Disjoint, _) => HealAdmissionResult::Accepted, + (OverlapVerdict::SameTarget, HealOverlapPolicy::Merge) => HealAdmissionResult::Merged, + (OverlapVerdict::SameTarget, _) => HealAdmissionResult::Dropped(HealAdmissionDropReason::AlreadyRunning), + (OverlapVerdict::Overlapping, _) => HealAdmissionResult::Dropped(HealAdmissionDropReason::OverlappingPaths), + }; + assert_eq!(receipt.result, expected, "{state:?} {policy:?} {existing} -> {incoming}"); + assert_eq!( + receipt.task_id, + if relation == OverlapVerdict::Disjoint { + request_id.clone() + } else { + owner + } + ); + if matches!(expected, HealAdmissionResult::Dropped(_)) { + assert!( + !manager + .heal_queue + .lock() + .await + .requests() + .any(|queued| queued.id == request_id) + ); + assert!(manager.task_aliases.lock().await.is_empty(), "rejection must not create a token alias"); + } + } + } + } +} + +#[tokio::test] +async fn admin_overlap_incompatible_options_conflict_without_replacing_settings() { + for state in [OwnerState::Active, OwnerState::Queued, OwnerState::Retrying] { + for option in ["dry_run", "scan_mode", "remove", "recreate", "parity", "recursive"] { + let manager = manager_with_policy(HealOverlapPolicy::Merge); + let original = admin_prefix_request("bucket", "scope/"); + let expected_options = original.options.clone(); + let owner = install_owner(&manager, original, state).await; + let mut request = admin_prefix_request("bucket", "scope/"); + match option { + "dry_run" => request.options.dry_run = true, + "scan_mode" => request.options.scan_mode = HealScanMode::Deep, + "remove" => request.options.remove_corrupted = true, + "recreate" => request.options.recreate_missing = false, + "parity" => request.options.update_parity = false, + "recursive" => request.options.recursive = true, + _ => unreachable!("fixture option"), + } + let receipt = manager + .submit_heal_request_with_receipt(request) + .await + .expect("option conflict"); + assert_eq!( + receipt.result, + HealAdmissionResult::Dropped(HealAdmissionDropReason::AlreadyRunning), + "{state:?} {option}" + ); + assert_eq!(receipt.task_id, owner); + assert_eq!( + manager + .get_task_report(&owner) + .await + .expect("original settings remain queryable") + .options, + Some(expected_options) + ); + } + } +} + +#[tokio::test] +async fn admin_overlap_consumed_timeout_does_not_break_equivalent_token_reuse() { + let manager = manager_with_policy(HealOverlapPolicy::Merge); + let mut original = admin_prefix_request("bucket", "scope/"); + original.options.timeout = Some(Duration::from_secs(17)); + let owner = install_owner(&manager, original, OwnerState::Retrying).await; + let receipt = manager + .submit_heal_request_with_receipt(admin_prefix_request("bucket", "scope/")) + .await + .expect("retry owner retains its token"); + assert_eq!(receipt.result, HealAdmissionResult::Merged); + assert_eq!(receipt.task_id, owner); + assert_eq!( + manager.retrying_heals.lock().await[&owner].request.options.timeout, + Some(Duration::from_secs(17)) + ); +} + +#[tokio::test] +async fn admin_overlap_typed_s3_targets_do_not_confuse_objects_and_prefixes() { + let object = |key: &str| admin_object_request(key).heal_type; + let prefix = |key: &str| admin_prefix_request("bucket", key).heal_type; + let erasure_set = HealType::ErasureSet { + buckets: vec![], + set_disk_id: "pool_0_set_1".to_string(), + }; + for (left, right, expected) in [ + (object("foo"), object("foobar"), OverlapVerdict::Disjoint), + (prefix("foo"), prefix("foobar"), OverlapVerdict::Overlapping), + (prefix("foo/"), object("foobar"), OverlapVerdict::Disjoint), + (object("foo"), prefix("foo/"), OverlapVerdict::Disjoint), + (prefix("scope/"), object("scope/child"), OverlapVerdict::Overlapping), + (object("/foo"), object("foo"), OverlapVerdict::Disjoint), + (object("foo/"), object("foo"), OverlapVerdict::Disjoint), + (prefix("scope%2F"), prefix("scope/"), OverlapVerdict::Disjoint), + (prefix("中文/"), object("中文/文件"), OverlapVerdict::Overlapping), + (HealType::Cluster, prefix("scope/"), OverlapVerdict::Overlapping), + ( + HealType::Bucket { + bucket: "bucket".to_string(), + }, + prefix("scope/"), + OverlapVerdict::Overlapping, + ), + (erasure_set.clone(), prefix("scope/"), OverlapVerdict::Overlapping), + ( + HealType::ErasureSet { + buckets: vec!["other".to_string()], + set_disk_id: "pool_0_set_1".to_string(), + }, + prefix("scope/"), + OverlapVerdict::Disjoint, + ), + ( + HealType::Object { + bucket: "bucket".to_string(), + object: "foo".to_string(), + version_id: Some("version-1".to_string()), + }, + HealType::Object { + bucket: "bucket".to_string(), + object: "foo".to_string(), + version_id: Some("version-2".to_string()), + }, + OverlapVerdict::Overlapping, + ), + ] { + for (existing, incoming) in [(&left, &right), (&right, &left)] { + let manager = manager_with_policy(HealOverlapPolicy::Merge); + let mut owner_request = HealRequest::new(existing.clone(), HealOptions::default(), HealPriority::Normal); + owner_request.source = HealRequestSource::Admin; + install_owner(&manager, owner_request, OwnerState::Active).await; + let mut request = HealRequest::new(incoming.clone(), HealOptions::default(), HealPriority::Normal); + request.source = HealRequestSource::Admin; + let receipt = manager + .submit_heal_request_with_receipt(request) + .await + .expect("typed target admission"); + assert_eq!( + receipt.result, + if expected == OverlapVerdict::Disjoint { + HealAdmissionResult::Accepted + } else { + HealAdmissionResult::Dropped(HealAdmissionDropReason::OverlappingPaths) + }, + "{existing:?} -> {incoming:?}" + ); + } + } +} + +#[tokio::test] +async fn admin_overlap_pool_set_scope_is_consistent_across_owner_states() { + for state in [OwnerState::Active, OwnerState::Queued, OwnerState::Retrying] { + for object in [false, true] { + for (incoming_pool, incoming_set, expected) in [ + (Some(0), Some(1), HealAdmissionResult::Merged), + (Some(0), Some(2), HealAdmissionResult::Accepted), + (Some(1), Some(1), HealAdmissionResult::Accepted), + (Some(0), None, HealAdmissionResult::Dropped(HealAdmissionDropReason::OverlappingPaths)), + (None, None, HealAdmissionResult::Dropped(HealAdmissionDropReason::OverlappingPaths)), + ] { + let make_request = || { + if object { + admin_object_request("scope/object") + } else { + admin_prefix_request("bucket", "scope/") + } + }; + let manager = manager_with_policy(HealOverlapPolicy::Merge); + let mut original = make_request(); + original.options.pool_index = Some(0); + original.options.set_index = Some(1); + let owner = install_owner(&manager, original, state).await; + let mut incoming = make_request(); + incoming.options.pool_index = incoming_pool; + incoming.options.set_index = incoming_set; + let id = incoming.id.clone(); + let receipt = manager + .submit_heal_request_with_receipt(incoming) + .await + .expect("scoped admission"); + assert_eq!( + receipt.result, expected, + "{state:?} object={object}, pool={incoming_pool:?}, set={incoming_set:?}" + ); + assert_eq!( + receipt.task_id, + if expected == HealAdmissionResult::Accepted { + id + } else { + owner + } + ); + } + } + } +} + +#[tokio::test] +async fn admin_overlap_background_admission_and_strict_policy_keep_their_boundaries() { + for policy in [HealOverlapPolicy::Merge, HealOverlapPolicy::MinioError] { + for background in [ + HealRequestSource::Scanner, + HealRequestSource::AutoHeal, + HealRequestSource::ReadRepair, + HealRequestSource::Internal, + ] { + let manager = manager_with_policy(policy); + install_owner(&manager, admin_prefix_request("bucket", "scope/"), OwnerState::Active).await; + let mut request = admin_prefix_request("bucket", "scope/child/"); + request.source = background; + assert_eq!( + manager + .submit_heal_request(request) + .await + .expect("background source remains admitted"), + HealAdmissionResult::Accepted + ); + + let manager = manager_with_policy(policy); + let mut original = admin_prefix_request("bucket", "scope/"); + original.source = background; + install_owner(&manager, original, OwnerState::Active).await; + assert_eq!( + manager + .submit_heal_request(admin_prefix_request("bucket", "scope/child/")) + .await + .expect("admin policy decision"), + if policy == HealOverlapPolicy::Merge { + HealAdmissionResult::Accepted + } else { + HealAdmissionResult::Dropped(HealAdmissionDropReason::OverlappingPaths) + } + ); + } + } +} + +#[tokio::test] +async fn admin_overlap_force_start_cancels_every_overlapping_state_and_keeps_disjoint_work() { + for policy in [HealOverlapPolicy::Merge, HealOverlapPolicy::MinioError] { + let manager = manager_with_policy(policy); + let active = install_owner(&manager, admin_prefix_request("bucket", "scope/active/"), OwnerState::Active).await; + let queued = install_owner(&manager, admin_prefix_request("bucket", "scope/queued/"), OwnerState::Queued).await; + let retrying = install_owner(&manager, admin_prefix_request("bucket", "scope/retrying/"), OwnerState::Retrying).await; + let retry_cancel = manager.retrying_heals.lock().await[&retrying].cancel_token.clone(); + let disjoint = install_owner(&manager, admin_prefix_request("bucket", "other/"), OwnerState::Queued).await; + let mut replacement = admin_prefix_request("bucket", "scope/"); + replacement.force_start = true; + let replacement_id = replacement.id.clone(); + let receipt = manager + .submit_heal_request_with_receipt(replacement) + .await + .expect("replace overlapping owners"); + assert_eq!(receipt.result, HealAdmissionResult::Accepted); + assert_eq!(receipt.task_id, replacement_id); + assert_eq!( + manager.get_task_status(&active).await.expect("active cancellation retained"), + HealTaskStatus::Cancelled + ); + assert!(retry_cancel.is_cancelled()); + assert!(!manager.retrying_heals.lock().await.contains_key(&retrying)); + let ids = manager + .heal_queue + .lock() + .await + .requests() + .map(|request| request.id.clone()) + .collect::>(); + assert_eq!(ids, HashSet::from([disjoint, replacement_id])); + assert!(!ids.contains(&queued)); + } +} + +#[tokio::test] +async fn admin_overlap_concurrent_parent_child_starts_admit_exactly_one_owner() { + for _ in 0..8 { + let manager = manager_with_policy(HealOverlapPolicy::Merge); + let (parent, child) = tokio::join!( + manager.submit_heal_request_with_receipt(admin_prefix_request("bucket", "scope/")), + manager.submit_heal_request_with_receipt(admin_prefix_request("bucket", "scope/child/")) + ); + let receipts = [parent.expect("parent decision"), child.expect("child decision")]; + assert_eq!( + receipts + .iter() + .filter(|receipt| receipt.result == HealAdmissionResult::Accepted) + .count(), + 1 + ); + assert_eq!( + receipts + .iter() + .filter(|receipt| receipt.result == HealAdmissionResult::Dropped(HealAdmissionDropReason::OverlappingPaths)) + .count(), + 1 + ); + assert_eq!(receipts[0].task_id, receipts[1].task_id); + assert_eq!(manager.heal_queue.lock().await.len(), 1); + } +} + +#[tokio::test] +async fn admin_overlap_normal_start_cannot_enter_force_start_cancellation_window() { + let manager = manager_with_policy(HealOverlapPolicy::Merge); + install_owner(&manager, admin_prefix_request("bucket", "scope/child/"), OwnerState::Active).await; + // Cancellation resolves aliases before acquiring runtime state. Holding + // this lock leaves the registry available inside the replacement window. + let aliases = manager.task_aliases.lock().await; + let mut replacement = admin_prefix_request("bucket", "scope/"); + replacement.force_start = true; + let mut forced = Box::pin(manager.submit_heal_request_with_receipt(replacement)); + tokio::time::timeout( + Duration::from_secs(5), + std::future::poll_fn(|cx| { + assert!( + std::pin::pin!(tokio::task::unconstrained(forced.as_mut())) + .poll(cx) + .is_pending() + ); + if manager.active_heals.try_lock().is_ok() { + Poll::Ready(()) + } else { + // Another test can briefly hold the shared admission probe lock. + // Drive the start until cancellation is blocked on our alias gate. + cx.waker().wake_by_ref(); + Poll::Pending + } + }), + ) + .await + .expect("forceStart reaches cancellation with registry locks released"); + let mut normal = Box::pin(manager.submit_heal_request_with_receipt(admin_prefix_request("bucket", "scope/other/"))); + assert!( + futures::poll!(tokio::task::unconstrained(normal.as_mut())).is_pending(), + "ordinary START must wait for the replacement gate" + ); + drop(aliases); + let (forced, normal) = tokio::join!(forced, normal); + assert_eq!(forced.expect("forced decision").result, HealAdmissionResult::Accepted); + assert_eq!( + normal.expect("normal decision").result, + HealAdmissionResult::Dropped(HealAdmissionDropReason::OverlappingPaths) + ); + assert_eq!(manager.heal_queue.lock().await.len(), 1); +} + +#[tokio::test] +async fn admin_overlap_cancel_cannot_miss_a_queued_to_active_transition() { + let manager = manager_with_policy(HealOverlapPolicy::Merge); + let owner = install_owner(&manager, admin_prefix_request("bucket", "scope/"), OwnerState::Queued).await; + let retrying = manager.retrying_heals.lock().await; + let mut cancelled = Box::pin(manager.cancel_task(&owner)); + assert!(futures::poll!(cancelled.as_mut()).is_pending()); + let mut scheduled = Box::pin(process_manager_queue_once(&manager)); + assert!( + futures::poll!(scheduled.as_mut()).is_pending(), + "scheduler cannot move the owner between cancellation lookups" + ); + drop(retrying); + let (cancelled, ()) = tokio::join!(cancelled, scheduled); + cancelled.expect("queued owner is cancelled"); + assert!(manager.heal_queue.lock().await.is_empty()); + assert!(!manager.active_heals.lock().await.contains_key(&owner)); +} + +#[tokio::test] +async fn admin_overlap_forced_request_id_replay_has_no_cancellation_side_effects() { + let manager = manager_with_policy(HealOverlapPolicy::Merge); + let mut original = admin_prefix_request("bucket", "scope/"); + original.force_start = true; + let original_id = install_owner(&manager, original.clone(), OwnerState::Queued).await; + // Older releases can leave independently accepted, overlapping owners. + let other_id = install_owner(&manager, admin_prefix_request("bucket", "scope/child/"), OwnerState::Active).await; + let receipt = manager + .submit_heal_request_with_receipt(original) + .await + .expect("replay original receipt"); + assert_eq!(receipt.result, HealAdmissionResult::Accepted); + assert_eq!(receipt.task_id, original_id); + assert!( + manager.active_heals.lock().await.contains_key(&other_id), + "receipt replay must not repeat cancellation" + ); + assert_eq!(manager.heal_queue.lock().await.len(), 1); +} + +struct RetryQueueHook { + reached: Notify, + release: Notify, + resumed: Notify, +} + +static RETRY_QUEUE_HOOKS: LazyLock>>> = + LazyLock::new(|| StdMutex::new(HashMap::new())); + +pub(in crate::heal::manager) async fn pause_before_retry_queue(task_id: &str) { + let hook = RETRY_QUEUE_HOOKS.lock().expect("retry queue hooks").get(task_id).cloned(); + if let Some(hook) = hook { + hook.reached.notify_one(); + hook.release.notified().await; + hook.resumed.notify_one(); + } +} + +#[tokio::test] +async fn admin_overlap_cancelled_retry_cannot_requeue_after_its_backoff_checks() { + let manager = manager_with_policy(HealOverlapPolicy::Merge); + let mut original = admin_object_request("object"); + original.heal_type = HealType::Object { + bucket: "retry-transition".to_string(), + object: "object".to_string(), + version_id: None, + }; + let owner = original.id.clone(); + let hook = Arc::new(RetryQueueHook { + reached: Notify::new(), + release: Notify::new(), + resumed: Notify::new(), + }); + RETRY_QUEUE_HOOKS + .lock() + .expect("install retry hook") + .insert(owner.clone(), Arc::clone(&hook)); + manager + .submit_heal_request(original) + .await + .expect("admit retryable object heal"); + process_manager_queue_once(&manager).await; + tokio::time::timeout(Duration::from_secs(10), hook.reached.notified()) + .await + .expect("real executor reaches retry queue acquisition"); + manager + .cancel_task(&owner) + .await + .expect("cancel after the retry ownership prechecks"); + hook.release.notify_one(); + tokio::time::timeout(Duration::from_secs(5), hook.resumed.notified()) + .await + .expect("retry worker resumes"); + // The test runs on a single-threaded runtime. Once this notification is + // observed, the worker has executed the uncontended queue ownership check. + assert!(manager.heal_queue.lock().await.is_empty(), "cancelled retry must not be republished"); + assert!(!manager.retrying_heals.lock().await.contains_key(&owner)); + RETRY_QUEUE_HOOKS.lock().expect("remove retry hook").remove(&owner); +} diff --git a/crates/heal/src/heal/manager/tests/root_recovery.rs b/crates/heal/src/heal/manager/tests/root_recovery.rs index 93cded92e..3dfc481d0 100644 --- a/crates/heal/src/heal/manager/tests/root_recovery.rs +++ b/crates/heal/src/heal/manager/tests/root_recovery.rs @@ -915,6 +915,161 @@ async fn root_recovery_force_start_cancels_only_overlapping_durable_admin_record assert_eq!(queued_ids, HashSet::from([disjoint.id, replacement.id])); } +#[tokio::test] +async fn root_recovery_admin_overlap_rejects_durable_only_owners_without_writing_new_intents() { + for (existing, incoming, reason) in [ + ("scope/", "scope/", HealAdmissionDropReason::AlreadyRunning), + ("scope/", "scope/child/", HealAdmissionDropReason::OverlappingPaths), + ("scope/child/", "scope/", HealAdmissionDropReason::OverlappingPaths), + ] { + let (_temp, disk) = recovery_disk().await; + let manager = recovery_manager(vec![disk]); + let mut owner = admin_prefix_request("bucket", existing); + owner.options.timeout = Some(Duration::ZERO); + manager.root_recovery.persist(&owner).await.expect("persist unreplayed owner"); + let receipt = manager + .submit_heal_request_with_receipt(admin_prefix_request("bucket", incoming)) + .await + .expect("durable overlap decision"); + assert_eq!(receipt.result, HealAdmissionResult::Dropped(reason)); + assert_eq!(receipt.task_id, owner.id); + assert!(manager.heal_queue.lock().await.is_empty()); + let pending = manager + .root_recovery + .pending() + .await + .expect("original durable responsibility remains"); + assert_eq!(pending.len(), 1); + assert_eq!(pending[0].id, owner.id); + assert_eq!( + pending[0].options.timeout, + Some(Duration::ZERO), + "admission must not reset an exhausted budget" + ); + } +} + +#[tokio::test] +async fn root_recovery_admin_overlap_same_id_does_not_overwrite_a_durable_only_budget() { + let (_temp, disk) = recovery_disk().await; + let manager = recovery_manager(vec![disk]); + let mut owner = admin_prefix_request("bucket", "scope/"); + owner.options.timeout = Some(Duration::ZERO); + manager.root_recovery.persist(&owner).await.expect("exhausted owner"); + let mut replay = owner.clone(); + replay.options.timeout = Some(Duration::from_secs(60)); + replay.force_start = true; + let receipt = manager + .submit_heal_request_with_receipt(replay) + .await + .expect("same ID conflicts with durable owner"); + assert_eq!(receipt.result, HealAdmissionResult::Dropped(HealAdmissionDropReason::AlreadyRunning)); + let pending = manager.root_recovery.pending().await.expect("retained owner"); + assert_eq!(pending.len(), 1); + assert_eq!(pending[0].options.timeout, Some(Duration::ZERO)); + assert!(manager.heal_queue.lock().await.is_empty()); +} + +#[tokio::test] +async fn root_recovery_admin_overlap_corrupt_preflight_does_not_cancel_a_live_owner() { + let (_temp, disk) = recovery_disk().await; + let manager = recovery_manager(vec![disk.clone()]); + let owner = admin_prefix_request("bucket", "scope/child/"); + manager + .submit_heal_request(owner.clone()) + .await + .expect("admit original owner"); + let corrupt = root_request(); + let path = format!("root-heal-{}.json", corrupt.id); + disk.write_all(RUSTFS_META_BUCKET, &path, b"{".to_vec().into()) + .await + .expect("inject corrupt ownership record"); + let mut replacement = admin_prefix_request("bucket", "scope/"); + replacement.force_start = true; + assert!( + manager.submit_heal_request(replacement).await.is_err(), + "unknown ownership must fail before cancellation" + ); + assert_eq!( + manager + .heal_queue + .lock() + .await + .requests() + .map(|request| request.id.clone()) + .collect::>(), + vec![owner.id.clone()] + ); + assert_eq!( + manager + .get_task_status(&owner.id) + .await + .expect("original owner remains queryable"), + HealTaskStatus::Pending + ); + assert_eq!( + disk.read_all(RUSTFS_META_BUCKET, &path) + .await + .expect("retain corrupt record") + .as_ref(), + b"{" + ); +} + +#[tokio::test] +async fn root_recovery_admin_overlap_preserves_legacy_owners_and_replays_replacement_cancellations() { + let (_temp, disk) = recovery_disk().await; + let manager = recovery_manager(vec![disk.clone()]); + let parent = admin_prefix_request("bucket", "scope/"); + let child = admin_prefix_request("bucket", "scope/child/"); + manager + .root_recovery + .persist(&parent) + .await + .expect("legacy parent responsibility"); + manager + .root_recovery + .persist(&child) + .await + .expect("legacy child responsibility"); + manager + .replay_root_heals() + .await + .expect("accepted legacy owners must not be discarded"); + assert_eq!(manager.heal_queue.lock().await.len(), 2); + let rejected = manager + .submit_heal_request_with_receipt(admin_prefix_request("bucket", "scope/child/deep/")) + .await + .expect("new overlap must reject after replay"); + assert_eq!(rejected.result, HealAdmissionResult::Dropped(HealAdmissionDropReason::OverlappingPaths)); + let mut replacement = admin_prefix_request("bucket", "scope/"); + replacement.force_start = true; + let receipt = manager + .submit_heal_request_with_receipt(replacement) + .await + .expect("replace both recovered owners"); + assert_eq!(receipt.result, HealAdmissionResult::Accepted); + drop(manager); + let restarted = recovery_manager(vec![disk]); + restarted.replay_root_heals().await.expect("restart replacement"); + assert_eq!( + restarted + .heal_queue + .lock() + .await + .requests() + .map(|request| request.id.clone()) + .collect::>(), + vec![receipt.task_id] + ); + for owner in [&parent.id, &child.id] { + assert_eq!( + restarted.get_task_status(owner).await.expect("cancellation survives restart"), + HealTaskStatus::Cancelled + ); + } +} + #[tokio::test] async fn root_recovery_queued_non_root_admin_owner_is_not_priority_displaced() { let (_temp, disk) = recovery_disk().await; @@ -1040,7 +1195,9 @@ async fn root_recovery_queued_owner_is_not_priority_displaced() { HealOptions::default(), HealPriority::Urgent, ); - bucket.source = HealRequestSource::Admin; + // Internal urgent work has the same displacement eligibility without + // taking the administrator overlap rejection before the capacity check. + bucket.source = HealRequestSource::Internal; assert_eq!( manager .submit_heal_request(bucket) diff --git a/crates/protos/src/heal_control.rs b/crates/protos/src/heal_control.rs index 1742500c5..b3d7c8282 100644 --- a/crates/protos/src/heal_control.rs +++ b/crates/protos/src/heal_control.rs @@ -353,10 +353,10 @@ pub enum Admission { Full, DroppedQueueFull, DroppedPolicy, - /// HS-06: admin start rejected because the same target is already being - /// healed (RUSTFS_HEAL_OVERLAP_POLICY=minio_error only). + /// Admin start rejected because the same target is already owned and + /// cannot be merged under the selected overlap policy. DroppedAlreadyRunning, - /// HS-06: admin start rejected because its path overlaps an active heal. + /// Admin start rejected because its scope overlaps an existing owner. DroppedOverlappingPaths, } diff --git a/docs/operations/scanner-runtime-controls.md b/docs/operations/scanner-runtime-controls.md index fd6f7ab6b..360de78eb 100644 --- a/docs/operations/scanner-runtime-controls.md +++ b/docs/operations/scanner-runtime-controls.md @@ -292,7 +292,7 @@ Heal knobs are environment-only and read by `HealConfig::default` (`crates/heal/ | `RUSTFS_HEAL_MAINLINE_READ_UTILIZATION_HIGH_PERCENT` | `80` (`DEFAULT_HEAL_MAINLINE_READ_UTILIZATION_HIGH_PERCENT`, capped at 100) | Read-utilization high watermark for start admission and running admin pacing; zero disables this class. | | `RUSTFS_HEAL_MAINLINE_WRITE_UTILIZATION_HIGH_PERCENT` | `80` (`DEFAULT_HEAL_MAINLINE_WRITE_UTILIZATION_HIGH_PERCENT`, capped at 100) | Write-utilization high watermark for start admission and running admin pacing; zero disables this class. | | `RUSTFS_HEAL_MAINLINE_MAX_SLEEP_MS` | `250` (`DEFAULT_HEAL_MAINLINE_MAX_SLEEP_MS`) | Start recheck interval; running admin waits cap each pacing-gate holder at 1000 ms. Zero disables running pacing. | -| `RUSTFS_HEAL_OVERLAP_POLICY` | `merge` (`DEFAULT_HEAL_OVERLAP_POLICY`) | `merge` dedups an admin heal start that overlaps a running or queued heal; `minio_error` returns a typed already-running / overlapping-paths rejection like madmin. | +| `RUSTFS_HEAL_OVERLAP_POLICY` | `merge` (`DEFAULT_HEAL_OVERLAP_POLICY`) | `merge` reuses the token of an equivalent active, queued, or retrying admin heal; incompatible same-target starts and intersecting admin scopes return `already_running` / `overlapping_paths`. Durable-only owners reject until recovered or replaced with `forceStart`. `minio_error` also rejects equivalent starts and preserves overlap rejection against background tasks. | | `RUSTFS_HEAL_MRF_ENABLE` | `true` (`DEFAULT_HEAL_MRF_ENABLE`) | MRF intent pipeline: error paths deliver repair intents to the heal runtime and unconsumed intents replay from the durable journal after restart. | | `RUSTFS_HEAL_MRF_QUEUE_SIZE` | `100000` (`DEFAULT_HEAL_MRF_QUEUE_SIZE`) | MRF in-memory queue capacity. | | `RUSTFS_HEAL_MRF_JOURNAL_MAX_BYTES` | `8388608` (`DEFAULT_HEAL_MRF_JOURNAL_MAX_BYTES`, 8 MiB) | MRF journal size at which compaction runs. | diff --git a/rustfs/src/admin/handlers/heal.rs b/rustfs/src/admin/handlers/heal.rs index 3be240b78..a33eda6ad 100644 --- a/rustfs/src/admin/handlers/heal.rs +++ b/rustfs/src/admin/handlers/heal.rs @@ -1989,6 +1989,7 @@ mod tests { ] { let error = reject_heal_admission(HealAdmissionResult::Dropped(reason)); assert_eq!(error.code(), &S3ErrorCode::OperationAborted); + assert_eq!(error.code().status_code(), Some(StatusCode::CONFLICT)); assert!( error.to_string().contains(label), "the caller must distinguish conflicts from transient coordination failure" diff --git a/rustfs/src/storage/rpc/node_service.rs b/rustfs/src/storage/rpc/node_service.rs index 4ba5f9461..41a383f11 100644 --- a/rustfs/src/storage/rpc/node_service.rs +++ b/rustfs/src/storage/rpc/node_service.rs @@ -3309,6 +3309,87 @@ mod tests { } if task_id == request_id)); } + #[tokio::test] + async fn heal_control_admin_overlap_receipts_preserve_token_and_conflict_reason() { + use rustfs_protos::heal_control::{Admission, Envelope, Outcome, RequestMetadata}; + + let (manager, mut parent, metadata) = heal_start_retry_fixture(); + parent.force_start = false; + parent.object_prefix = Some("scope/".to_string()); + let parent_id = parent.id.clone(); + let envelope = Envelope::start(parent.clone(), metadata).expect("parent start"); + let response = execute_heal_control_envelope_with_manager(envelope, metadata.coordinator_epoch, Some(manager.clone())) + .await + .expect("admit parent scope"); + assert!( + matches!(decode_transport_start_outcome(&response, &parent_id, metadata.coordinator_epoch), + Outcome::Start { task_id, admission: Admission::Accepted } if task_id == parent_id) + ); + + for (prefix, expected) in [ + ("scope/", Admission::Merged), + ("scope/child/", Admission::DroppedOverlappingPaths), + ("other/", Admission::Accepted), + ] { + let mut request = parent.clone(); + request.id = Uuid::new_v4().to_string(); + request.object_prefix = Some(prefix.to_string()); + let request_id = request.id.clone(); + let request_metadata = RequestMetadata { + nonce: *Uuid::new_v4().as_bytes(), + ..metadata + }; + let envelope = Envelope::start(request, request_metadata).expect("scoped start envelope"); + let response = + execute_heal_control_envelope_with_manager(envelope, metadata.coordinator_epoch, Some(manager.clone())) + .await + .expect("overlap remains a typed admission result across the RPC boundary"); + let Outcome::Start { task_id, admission } = + decode_transport_start_outcome(&response, &request_id, metadata.coordinator_epoch) + else { + panic!("start must return an admission receipt") + }; + assert_eq!(admission, expected, "prefix={prefix}"); + assert_eq!( + task_id, + if expected == Admission::Accepted { + request_id + } else { + parent_id.clone() + } + ); + } + assert_eq!(manager.operations_snapshot().await.queue_length, 2); + assert_eq!( + manager + .get_task_status(&parent_id) + .await + .expect("rejected child preserves parent"), + rustfs_heal::heal::task::HealTaskStatus::Pending + ); + + parent.id = Uuid::new_v4().to_string(); + parent.force_start = true; + let replacement_id = parent.id.clone(); + let envelope = Envelope::start( + parent, + RequestMetadata { + nonce: *Uuid::new_v4().as_bytes(), + ..metadata + }, + ) + .expect("force replacement envelope"); + let response = execute_heal_control_envelope_with_manager(envelope, metadata.coordinator_epoch, Some(manager.clone())) + .await + .expect("forceStart replaces the parent and preserves the disjoint scope"); + assert!( + matches!(decode_transport_start_outcome(&response, &replacement_id, metadata.coordinator_epoch), + Outcome::Start { task_id, admission: Admission::Accepted } if task_id == replacement_id) + ); + assert_ne!(replacement_id, parent_id); + assert_eq!(manager.operations_snapshot().await.queue_length, 2); + } + #[tokio::test] async fn heal_start_retry_new_forced_request_is_a_distinct_start() { let (manager, request, metadata) = heal_start_retry_fixture();