fix(heal): reject overlapping admin heal owners (#7723)

This commit is contained in:
cxymds
2026-09-13 10:55:35 +08:00
committed by GitHub
parent 7776004977
commit 2cfefcfd65
12 changed files with 1082 additions and 197 deletions
+3 -3
View File
@@ -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";
+4 -5
View File
@@ -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,
}
+248 -174
View File
@@ -515,7 +515,7 @@ fn publish_heal_queue_length(queue: &PriorityHealQueue) {
fn active_heal_for_dedup_key(active_heals: &HashMap<String, Arc<HealTask>>, 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<usize>, Option<usize>) {
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<HealAdmissionDropReason> {
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<std::sync::Mutex<HashSet<String>>>,
/// Durable handoff of interrupted administrator root traversals.
root_recovery: Arc<root_recovery::RootHealRecovery>,
/// 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<dyn HealStorageAPI>,
/// 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<String> = {
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::<Vec<_>>();
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::<HashSet<_>>();
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::<Vec<_>>();
let live_ids = live_owners.iter().map(|(id, _, _, _)| *id).collect::<HashSet<_>>();
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 {
+15 -4
View File
@@ -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}"),
}
}
+13 -2
View File
@@ -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(
+4 -4
View File
@@ -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);
@@ -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::<HashSet<_>>();
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<StdMutex<HashMap<String, Arc<RetryQueueHook>>>> =
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);
}
@@ -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<_>>(),
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<_>>(),
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)
+3 -3
View File
@@ -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,
}
+1 -1
View File
@@ -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. |
+1
View File
@@ -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"
+81
View File
@@ -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();