diff --git a/crates/ecstore/src/cluster/rpc/remote_disk.rs b/crates/ecstore/src/cluster/rpc/remote_disk.rs index efb1b4d8d..32bf3ae57 100644 --- a/crates/ecstore/src/cluster/rpc/remote_disk.rs +++ b/crates/ecstore/src/cluster/rpc/remote_disk.rs @@ -50,10 +50,10 @@ use rustfs_protos::evict_failed_connection; use rustfs_protos::proto_gen::node_service::RenamePartRequest; use rustfs_protos::proto_gen::node_service::{ BatchReadVersionRequest, BatchReadVersionResponse, CheckPartsRequest, DeletePathsRequest, DeleteRequest, - DeleteVersionRequest, DeleteVersionsRequest, DeleteVolumeRequest, DiskInfoRequest, ListDirRequest, ListVolumesRequest, - MakeVolumeRequest, MakeVolumesRequest, PreparePartTransactionRequest, ReadAllRequest, ReadMetadataRequest, - ReadMultipleRequest, ReadMultipleResponse, ReadPartsRequest, ReadVersionRequest, ReadXlRequest, RenameDataRequest, - RenameFileRequest, SettlePartTransactionRequest, SnapshotLeaseReleaseRequest, SnapshotLeaseRenewRequest, + DeleteVersionRequest, DeleteVersionsRequest, DeleteVersionsResponse, DeleteVolumeRequest, DiskInfoRequest, ListDirRequest, + ListVolumesRequest, MakeVolumeRequest, MakeVolumesRequest, PreparePartTransactionRequest, ReadAllRequest, + ReadMetadataRequest, ReadMultipleRequest, ReadMultipleResponse, ReadPartsRequest, ReadVersionRequest, ReadXlRequest, + RenameDataRequest, RenameFileRequest, SettlePartTransactionRequest, SnapshotLeaseReleaseRequest, SnapshotLeaseRenewRequest, SnapshotLeaseRequest, SnapshotLeaseResponse, StatVolumeRequest, UpdateMetadataRequest, VerifyFileRequest, WriteAllRequest, WriteMetadataRequest, node_service_client::NodeServiceClient, }; @@ -112,6 +112,28 @@ const EVENT_REMOTE_DISK_RPC: &str = "remote_disk_rpc"; const SNAPSHOT_LEASE_PROTOCOL_VERSION: u32 = 1; pub const REMOTE_SNAPSHOT_LEASE_TTL: Duration = Duration::from_secs(60); +fn decode_delete_versions_errors(response: DeleteVersionsResponse, expected_len: usize) -> Vec> { + if !response.item_errors.is_empty() { + if response.item_errors.len() != expected_len { + return vec![Some(Error::other("malformed delete_versions item errors")); expected_len]; + } + return response + .item_errors + .into_iter() + .map(|error| (error.code != 0).then(|| error.into())) + .collect(); + } + + if response.errors.len() != expected_len { + return vec![Some(Error::other("malformed delete_versions errors")); expected_len]; + } + response + .errors + .into_iter() + .map(|error| (!error.is_empty()).then(|| Error::other(error))) + .collect() +} + fn snapshot_lease_token_from_response(response: SnapshotLeaseResponse) -> Result { if !response.success { return Err(response.error.unwrap_or_default().into()); @@ -2406,8 +2428,6 @@ impl DiskAPI for RemoteDisk { return errors; } - // TODO(backlog): replace string errors with typed `StorageError` variants - let result = self .execute_with_timeout( || async { @@ -2439,17 +2459,7 @@ impl DiskAPI for RemoteDisk { } return errors; } - response - .errors - .iter() - .map(|error| { - if error.is_empty() { - None - } else { - Some(Error::other(error.to_string())) - } - }) - .collect() + decode_delete_versions_errors(response, versions.len()) } #[tracing::instrument(level = "trace", skip_all)] @@ -3760,6 +3770,63 @@ mod tests { static INIT: Once = Once::new(); + #[test] + fn delete_versions_response_preserves_typed_item_errors() { + let errors = decode_delete_versions_errors( + DeleteVersionsResponse { + success: true, + errors: vec!["file not found".to_string(), String::new()], + error: None, + item_errors: vec![ + rustfs_protos::proto_gen::node_service::Error { + code: DiskError::FileNotFound.to_u32(), + error_info: "file not found".to_string(), + }, + rustfs_protos::proto_gen::node_service::Error::default(), + ], + }, + 2, + ); + + assert!(matches!(errors.as_slice(), [Some(DiskError::FileNotFound), None])); + } + + #[test] + fn delete_versions_response_accepts_legacy_string_errors() { + let errors = decode_delete_versions_errors( + DeleteVersionsResponse { + success: true, + errors: vec!["legacy error".to_string(), String::new()], + error: None, + item_errors: Vec::new(), + }, + 2, + ); + + assert_eq!(errors.len(), 2); + assert_eq!(errors[0].as_ref().map(ToString::to_string).as_deref(), Some("io error legacy error")); + assert!(errors[1].is_none()); + } + + #[test] + fn delete_versions_response_rejects_misaligned_item_errors() { + let errors = decode_delete_versions_errors( + DeleteVersionsResponse { + success: true, + errors: vec!["file not found".to_string()], + error: None, + item_errors: vec![rustfs_protos::proto_gen::node_service::Error { + code: DiskError::FileNotFound.to_u32(), + error_info: "file not found".to_string(), + }], + }, + 2, + ); + + assert_eq!(errors.len(), 2); + assert!(errors.iter().all(Option::is_some)); + } + #[test] fn disk_mutation_digest_marks_rolling_compatibility() { let mut request = Request::new(()); diff --git a/crates/heal/src/heal/channel.rs b/crates/heal/src/heal/channel.rs index edfcf2613..e00435100 100644 --- a/crates/heal/src/heal/channel.rs +++ b/crates/heal/src/heal/channel.rs @@ -1640,6 +1640,60 @@ mod tests { assert_eq!(payload["items"].as_array().expect("items should be an array").len(), 0); } + #[tokio::test] + async fn test_process_query_request_reports_displaced_terminal_detail() { + let heal_manager = Arc::new(HealManager::new( + Arc::new(MockStorage), + Some(HealConfig { + queue_size: 1, + ..HealConfig::default() + }), + )); + let mut displaced = HealRequest::new( + HealType::Bucket { + bucket: "displaced-channel".to_string(), + }, + HealOptions::default(), + HealPriority::Low, + ); + displaced.id = "displaced-channel-task".to_string(); + let displaced_id = displaced.id.clone(); + heal_manager + .submit_heal_request(displaced) + .await + .expect("initial channel task should queue"); + heal_manager + .submit_heal_request(HealRequest::new( + HealType::Bucket { + bucket: "successor-channel".to_string(), + }, + HealOptions::default(), + HealPriority::High, + )) + .await + .expect("successor channel task should displace the initial task"); + + let processor = HealChannelProcessor::new(heal_manager); + let (tx, rx) = oneshot::channel(); + processor + .process_query_request("displaced-channel".to_string(), displaced_id, None, tx) + .await + .expect("displaced query should process"); + let response = rx + .await + .expect("query response should be returned") + .expect("displaced query should remain successful"); + let payload: serde_json::Value = serde_json::from_slice(response.data.as_deref().expect("status payload should exist")) + .expect("status payload should be json"); + assert_eq!(payload["summary"], "stopped"); + assert!( + response + .error + .as_deref() + .is_some_and(|detail| detail.contains("reason=displaced")) + ); + } + #[tokio::test] async fn test_process_query_request_reports_running_for_queued_task() { let heal_manager = create_test_heal_manager(); diff --git a/crates/heal/src/heal/manager.rs b/crates/heal/src/heal/manager.rs index bad6d4b5b..fc8255e0b 100644 --- a/crates/heal/src/heal/manager.rs +++ b/crates/heal/src/heal/manager.rs @@ -40,6 +40,7 @@ use tracing::{debug, error, info, warn}; use super::{DiskError, Endpoint, HealDiskExt as _, local_disk_map_read}; const KEEP_HEAL_TASK_STATUS_DURATION: Duration = Duration::from_secs(10 * 60); +const DISPLACED_HEAL_REASON: &str = "reason=displaced; retry_hint=submit_again"; const LOG_COMPONENT_HEAL: &str = "heal"; const LOG_SUBSYSTEM_DISK_SCANNER: &str = "disk_scanner"; const LOG_SUBSYSTEM_MANAGER: &str = "manager"; @@ -120,26 +121,30 @@ struct MrfRepairNoticeTarget { version_id: Option<[u8; 16]>, } -#[derive(Debug, Clone, PartialEq, Eq)] +#[derive(Debug, Clone)] struct HealAdmissionDecision { result: HealAdmissionResult, - displaced_task_id: Option, + displaced_request: Option, } impl HealAdmissionDecision { const fn new(result: HealAdmissionResult) -> Self { Self { result, - displaced_task_id: None, + displaced_request: None, } } - fn accepted_with_displacement(displaced_task_id: String) -> Self { + fn accepted_with_displacement(displaced_request: HealRequest) -> Self { Self { result: HealAdmissionResult::Accepted, - displaced_task_id: Some(displaced_task_id), + displaced_request: Some(displaced_request), } } + + fn displaced_task_id(&self) -> Option<&str> { + self.displaced_request.as_ref().map(|request| request.id.as_str()) + } } fn lock_mrf_repair_notice_targets( @@ -151,6 +156,55 @@ fn lock_mrf_repair_notice_targets( } } +fn lock_displaced_terminals( + registry: &StdMutex>>, +) -> StdMutexGuard<'_, HashMap>> { + match registry.lock() { + Ok(guard) => guard, + Err(poisoned) => poisoned.into_inner(), + } +} + +fn record_displaced_terminal( + registry: &StdMutex>>, + request: &HealRequest, +) -> Arc { + let terminal = Arc::new(CompletedHealStatus { + heal_type: request.heal_type.clone(), + status: HealTaskStatus::Failed { + error: format!("heal task displaced by a higher-priority request ({DISPLACED_HEAL_REASON})"), + }, + result_items_truncated: false, + completed_at: SystemTime::now(), + seqed_items: Vec::new(), + next_seq: 0, + min_seq: 0, + }); + let mut terminals = lock_displaced_terminals(registry); + prune_completed_heal_statuses(&mut terminals); + terminals.insert(request.id.clone(), Arc::clone(&terminal)); + terminal +} + +async fn remove_displaced_task_aliases( + aliases: &Arc>>, + terminals: &StdMutex>>, + task_id: &str, + terminal: &Arc, +) { + let mut aliases = aliases.lock().await; + let alias_ids = aliases + .iter() + .filter_map(|(alias_id, alias)| (alias.task_id == task_id).then_some(alias_id.clone())) + .collect::>(); + let mut displaced_terminals = lock_displaced_terminals(terminals); + prune_completed_heal_statuses(&mut displaced_terminals); + for alias_id in alias_ids { + displaced_terminals.insert(alias_id, Arc::clone(terminal)); + } + aliases.retain(|alias_id, alias| alias_id != task_id && alias.task_id != task_id); +} + async fn remove_task_aliases_for_task(registry: &Arc>>, task_id: &str) { registry .lock() @@ -618,6 +672,14 @@ pub struct HealManager { /// are shared so the lookup helper can hand a completed entry to a /// caller without cloning the retained result window. completed_heals: Arc>>>, + /// Terminals for requests removed by priority displacement. An Accepted + /// task ID remains queryable for the same process lifetime and the normal + /// ten-minute status TTL; clients should treat `reason=displaced` as a + /// terminal result and submit a fresh request. This sidecar is synchronous + /// so admission can publish the terminal while the queue transition is + /// still under its lock, without awaiting another tokio lock. Queue state + /// is process-local, so this guarantee does not extend across restart. + displaced_terminals: Arc>>>, /// Client tokens merged into an existing task id. task_aliases: Arc>>, /// Heal tasks waiting for a retry backoff to expire. @@ -659,6 +721,7 @@ struct HealQueueContext<'a> { heal_queue: &'a Arc>, active_heals: &'a Arc>>>, completed_heals: &'a Arc>>>, + displaced_terminals: &'a Arc>>>, task_aliases: &'a Arc>>, retrying_heals: &'a Arc>>, mrf_repair_notice_targets: &'a Arc>>>, @@ -874,7 +937,7 @@ impl HealManager { result = "accepted_by_displacement", "Heal queue request accepted by displacement" }); - return HealAdmissionDecision::accepted_with_displacement(displaced.id); + return HealAdmissionDecision::accepted_with_displacement(displaced); } demote_to_debug_when!(per_object_request, warn, target: "rustfs::heal::manager", { @@ -1105,6 +1168,7 @@ impl HealManager { active_heals: Arc::new(Mutex::new(HashMap::new())), heal_queue: Arc::new(Mutex::new(PriorityHealQueue::new())), completed_heals: Arc::new(Mutex::new(HashMap::new())), + displaced_terminals: Arc::new(StdMutex::new(HashMap::new())), task_aliases: Arc::new(Mutex::new(HashMap::new())), retrying_heals: Arc::new(Mutex::new(HashMap::new())), mrf_repair_notice_targets: Arc::new(StdMutex::new(HashMap::new())), @@ -1209,6 +1273,10 @@ impl HealManager { active_heals.clear(); publish_active_heal_count(&active_heals); self.completed_heals.lock().await.clear(); + // Do not let the synchronous guard live across the following async lock. + { + lock_displaced_terminals(&self.displaced_terminals).clear(); + } self.task_aliases.lock().await.clear(); self.retrying_heals.lock().await.clear(); lock_mrf_repair_notice_targets(&self.mrf_repair_notice_targets).clear(); @@ -1459,7 +1527,11 @@ impl HealManager { task_id = queued_id.to_owned(); } let should_notify = matches!(admission, HealAdmissionResult::Accepted) && config.event_driven_scheduler_enable; - let displaced_task_id = admission_decision.displaced_task_id; + let displaced_task_id = admission_decision.displaced_task_id().map(ToOwned::to_owned); + let displaced_terminal = admission_decision + .displaced_request + .as_ref() + .map(|request| record_displaced_terminal(&self.displaced_terminals, request)); if matches!(admission, HealAdmissionResult::Accepted | HealAdmissionResult::Merged) && let Some(target) = mrf_notice_target { @@ -1473,8 +1545,12 @@ impl HealManager { drop(queue); drop(active_heals); - if let Some(displaced_task_id) = displaced_task_id { - self.remove_aliases_for_task(&displaced_task_id).await; + if let (Some(displaced_task_id), Some(displaced_terminal)) = (displaced_task_id, displaced_terminal) { + // The queue has already removed the displaced request, so the + // synchronous terminal sidecar was published before aliases and + // MRF ownership are cleaned up. + remove_displaced_task_aliases(&self.task_aliases, &self.displaced_terminals, &displaced_task_id, &displaced_terminal) + .await; } if should_notify { @@ -1549,6 +1625,15 @@ impl HealManager { } } + if terminal_completed.is_none() { + let mut displaced_terminals = lock_displaced_terminals(&self.displaced_terminals); + prune_completed_heal_statuses(&mut displaced_terminals); + terminal_completed = displaced_terminals + .get(canonical_task_id) + .filter(|terminal| matches_path(&terminal.heal_type)) + .cloned(); + } + match terminal_completed { Some(completed) => TaskStateLookup::Completed(completed), None => TaskStateLookup::NotFound, @@ -1669,9 +1754,19 @@ impl HealManager { let mut completed_heals = self.completed_heals.lock().await; prune_completed_heal_statuses(&mut completed_heals); - completed_heals + if completed_heals .values() .any(|completed| heal_type_matches_path(&completed.heal_type, heal_path)) + { + return true; + } + drop(completed_heals); + + let mut displaced_terminals = lock_displaced_terminals(&self.displaced_terminals); + prune_completed_heal_statuses(&mut displaced_terminals); + displaced_terminals + .values() + .any(|terminal| heal_type_matches_path(&terminal.heal_type, heal_path)) } /// Get task progress diff --git a/crates/heal/src/heal/manager/auto_scan.rs b/crates/heal/src/heal/manager/auto_scan.rs index 10ac136a4..f7eb4482d 100644 --- a/crates/heal/src/heal/manager/auto_scan.rs +++ b/crates/heal/src/heal/manager/auto_scan.rs @@ -21,6 +21,7 @@ impl HealManager { let heal_queue = self.heal_queue.clone(); let active_heals = self.active_heals.clone(); let task_aliases = self.task_aliases.clone(); + let displaced_terminals = self.displaced_terminals.clone(); let mrf_repair_notice_targets = self.mrf_repair_notice_targets.clone(); let storage = self.storage.clone(); let replacement_recovery_anchors = self.replacement_recovery_anchors.clone(); @@ -481,6 +482,10 @@ impl HealManager { let admission = admission_decision.result; let should_notify = matches!(admission, HealAdmissionResult::Accepted) && config.event_driven_scheduler_enable; + let displaced_terminal = admission_decision + .displaced_request + .as_ref() + .map(|request| record_displaced_terminal(&displaced_terminals, request)); if matches!(admission, HealAdmissionResult::Accepted) && let Some(anchor) = recovery_anchor { @@ -491,8 +496,16 @@ impl HealManager { } drop(queue); drop(config); - if let Some(displaced_task_id) = admission_decision.displaced_task_id { - remove_task_aliases_for_task(&task_aliases, &displaced_task_id).await; + if let (Some(displaced_task_id), Some(displaced_terminal)) = + (admission_decision.displaced_task_id().map(ToOwned::to_owned), displaced_terminal) + { + remove_displaced_task_aliases( + &task_aliases, + &displaced_terminals, + &displaced_task_id, + &displaced_terminal, + ) + .await; lock_mrf_repair_notice_targets(&mrf_repair_notice_targets).remove(&displaced_task_id); } if matches!(admission, HealAdmissionResult::Accepted) { diff --git a/crates/heal/src/heal/manager/scheduler.rs b/crates/heal/src/heal/manager/scheduler.rs index d78330a9d..e471444d1 100644 --- a/crates/heal/src/heal/manager/scheduler.rs +++ b/crates/heal/src/heal/manager/scheduler.rs @@ -365,6 +365,7 @@ impl HealManager { let heal_queue = self.heal_queue.clone(); let active_heals = self.active_heals.clone(); let completed_heals = self.completed_heals.clone(); + let displaced_terminals = self.displaced_terminals.clone(); let task_aliases = self.task_aliases.clone(); let retrying_heals = self.retrying_heals.clone(); let mrf_repair_notice_targets = self.mrf_repair_notice_targets.clone(); @@ -397,6 +398,7 @@ impl HealManager { heal_queue: &heal_queue, active_heals: &active_heals, completed_heals: &completed_heals, + displaced_terminals: &displaced_terminals, task_aliases: &task_aliases, retrying_heals: &retrying_heals, mrf_repair_notice_targets: &mrf_repair_notice_targets, @@ -415,6 +417,7 @@ impl HealManager { heal_queue: &heal_queue, active_heals: &active_heals, completed_heals: &completed_heals, + displaced_terminals: &displaced_terminals, task_aliases: &task_aliases, retrying_heals: &retrying_heals, mrf_repair_notice_targets: &mrf_repair_notice_targets, @@ -442,6 +445,7 @@ impl HealManager { heal_queue, active_heals, completed_heals, + displaced_terminals, task_aliases, retrying_heals, mrf_repair_notice_targets, @@ -527,6 +531,7 @@ impl HealManager { let active_heals_clone = active_heals.clone(); let heal_queue_clone = heal_queue.clone(); let completed_heals_clone = completed_heals.clone(); + let displaced_terminals_clone = displaced_terminals.clone(); let task_aliases_clone = task_aliases.clone(); let retrying_heals_clone = retrying_heals.clone(); let mrf_repair_notice_targets_clone = mrf_repair_notice_targets.clone(); @@ -581,11 +586,167 @@ impl HealManager { "Heal scheduler task completed" ); } - Err(e) => { - let will_retry = retry_request.is_some(); - if will_retry { - demote_to_debug_when!(task.heal_type.is_per_object(), warn, target: "rustfs::heal::manager", { - event = EVENT_HEAL_SCHEDULER_STATE, + } + } + let retry_request_for_status = retry_request.as_ref().map(|(request, _, error)| HealTaskStatus::Retrying { + error: error.clone(), + retry_attempt: request.retry_attempts, + }); + let retry_request_for_queue = retry_request; + let retry_cancel_token = retry_request_for_queue.as_ref().map(|_| CancellationToken::new()); + if retry_request_for_queue.is_none() { + replacement_recovery_anchors_clone + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()) + .remove(&task_id); + } + let mut active_heals_guard = active_heals_clone.lock().await; + // Keep retry ownership continuous: status snapshots acquire + // these locks in the same active -> retrying order. + let mut retrying_heals_guard = if let (Some((request, _, error)), Some(cancel_token)) = + (retry_request_for_queue.as_ref(), retry_cancel_token.as_ref()) + { + let mut retrying = retrying_heals_clone.lock().await; + if active_heals_guard.contains_key(&task_id) { + retrying.insert( + request.id.clone(), + RetryingHeal { + request: request.clone(), + error: error.clone(), + cancel_token: cancel_token.clone(), + }, + ); + #[cfg(test)] + pause_retry_ownership_transition(&task_id, false).await; + } + Some(retrying) + } else { + None + }; + let completed_task = active_heals_guard.remove(&task_id); + if let Some(completed_task) = completed_task.as_ref() { + publish_active_heal_count(&active_heals_guard); + update_task_running_metric_for_task(&active_heals_guard, completed_task.as_ref()); + } + let active_count = active_heals_guard.len(); + drop(retrying_heals_guard.take()); + drop(active_heals_guard); + + if let Some(completed_task) = completed_task { + let completed_status = if let Some(status) = retry_request_for_status { + status + } else { + completed_task.get_status().await + }; + let terminal_completion = !matches!(completed_status, HealTaskStatus::Retrying { .. }); + let successful_completion = matches!(completed_status, HealTaskStatus::Completed); + let completed_progress = completed_task.get_progress().await; + // Single snapshot of the retained window: the task is + // finished and already off the active map, so there is + // no concurrent writer to race with. + let seqed_items = completed_task.get_seqed_result_items().await; + let (next_seq, min_seq) = completed_task.result_seq_cursors(); + let completed_status_entry = CompletedHealStatus { + heal_type: completed_task.heal_type.clone(), + status: completed_status.clone(), + result_items_truncated: completed_task.result_items_truncated(), + completed_at: SystemTime::now(), + seqed_items, + next_seq, + min_seq, + }; + let mut completed_heals_guard = completed_heals_clone.lock().await; + prune_completed_heal_statuses(&mut completed_heals_guard); + completed_heals_guard.insert(task_id.clone(), Arc::new(completed_status_entry)); + drop(completed_heals_guard); + // update statistics + let mut stats = statistics_clone.write().await; + match completed_status { + HealTaskStatus::Completed => { + stats.update_task_completion(true); + stats.add_healed_objects(completed_progress.objects_healed, completed_progress.bytes_processed); + } + HealTaskStatus::Retrying { .. } => {} + _ => { + stats.update_task_completion(false); + } + } + stats.update_running_tasks(usize_to_u64_saturated(active_count)); + drop(stats); + if terminal_completion { + let notice_targets = take_mrf_repair_notice_targets(&mrf_repair_notice_targets_clone, &task_id); + if successful_completion { + emit_mrf_repaired_events(notice_targets); + } + task_aliases_clone + .lock() + .await + .retain(|alias_id, alias| alias_id != &task_id && alias.task_id != task_id); + } + } + + if let (Some((retry_request, retry_delay, retry_error)), Some(retry_cancel_token)) = + (retry_request_for_queue, retry_cancel_token) + { + let retry_request_id = retry_request.id.clone(); + let retry_attempt = retry_request.retry_attempts; + let retry_key = PriorityHealQueue::make_dedup_key(&retry_request); + let retry_priority = retry_request.priority; + let retry_active_heals = active_heals_clone.clone(); + let retry_heal_queue = heal_queue_clone.clone(); + let retrying_heals_for_spawn = retrying_heals_clone.clone(); + let retry_task_aliases = task_aliases_clone.clone(); + let retry_displaced_terminals = displaced_terminals_clone.clone(); + let retry_mrf_repair_notice_targets = mrf_repair_notice_targets_clone.clone(); + let retry_completed_heals = completed_heals_clone.clone(); + let retry_notify = notify_clone.clone(); + let retry_manager_cancel_token = manager_cancel_token.clone(); + let retry_config = config_for_spawn.clone(); + tokio::spawn(async move { + loop { + tokio::select! { + _ = retry_cancel_token.cancelled() => { + debug!( + target: "rustfs::heal::manager", + event = EVENT_HEAL_QUEUE_ADMISSION, + component = LOG_COMPONENT_HEAL, + subsystem = LOG_SUBSYSTEM_MANAGER, + request_id = %retry_request_id, + priority = ?retry_priority, + retry_attempt, + result = "retry_cancelled", + "Heal retry admission decided" + ); + return; + } + _ = retry_manager_cancel_token.cancelled() => { + retrying_heals_for_spawn.lock().await.remove(&retry_request_id); + return; + } + _ = sleep(retry_delay) => {} + } + + { + let retrying_heals_guard = retrying_heals_for_spawn.lock().await; + if !retrying_heals_guard.contains_key(&retry_request_id) { + return; + } + } + + let active_duplicate_task_id = { + let active_heals_guard = retry_active_heals.lock().await; + active_heal_for_dedup_key(&active_heals_guard, &retry_key).map(|(task_id, _)| task_id) + }; + if let Some(active_duplicate_task_id) = active_duplicate_task_id { + retrying_heals_for_spawn.lock().await.remove(&retry_request_id); + move_mrf_repair_notice_targets( + &retry_mrf_repair_notice_targets, + &retry_request_id, + &active_duplicate_task_id, + ); + debug!( + target: "rustfs::heal::manager", + event = EVENT_HEAL_QUEUE_ADMISSION, component = LOG_COMPONENT_HEAL, subsystem = LOG_SUBSYSTEM_MANAGER, task_id, @@ -657,120 +818,43 @@ impl HealManager { drop(retrying_heals_guard.take()); drop(active_heals_guard); - if let Some(completed_task) = completed_task { - let completed_status = if let Some(status) = retry_request_for_status { - status - } else { - completed_task.get_status().await - }; - let terminal_completion = !matches!(completed_status, HealTaskStatus::Retrying { .. }); - let successful_completion = matches!(completed_status, HealTaskStatus::Completed); - let completed_progress = completed_task.get_progress().await; - // Single snapshot of the retained window: the task is - // finished and already off the active map, so there is - // no concurrent writer to race with. - let seqed_items = completed_task.get_seqed_result_items().await; - let (next_seq, min_seq) = completed_task.result_seq_cursors(); - let completed_status_entry = CompletedHealStatus { - heal_type: completed_task.heal_type.clone(), - status: completed_status.clone(), - result_items_truncated: completed_task.result_items_truncated(), - completed_at: SystemTime::now(), - seqed_items, - next_seq, - min_seq, - }; - let mut completed_heals_guard = completed_heals_clone.lock().await; - prune_completed_heal_statuses(&mut completed_heals_guard); - completed_heals_guard.insert(task_id.clone(), Arc::new(completed_status_entry)); - drop(completed_heals_guard); - // update statistics - let mut stats = statistics_clone.write().await; - match completed_status { - HealTaskStatus::Completed => { - stats.update_task_completion(true); - stats.add_healed_objects( - completed_progress.objects_healed, - completed_progress.bytes_processed, - ); - } - HealTaskStatus::Retrying { .. } => {} - _ => { - stats.update_task_completion(false); - } - } - stats.update_running_tasks(usize_to_u64_saturated(active_count)); - drop(stats); - #[cfg(test)] - panic_if_armed(SchedulerPanicPoint::Cleanup, &task_id); - if terminal_completion { - let notice_targets = take_mrf_repair_notice_targets(&mrf_repair_notice_targets_clone, &task_id); - if successful_completion { - emit_mrf_repaired_events(notice_targets); - } - task_aliases_clone - .lock() - .await - .retain(|alias_id, alias| alias_id != &task_id && alias.task_id != task_id); - } - } - - if let (Some((retry_request, retry_delay, retry_error)), Some(retry_cancel_token)) = - (retry_request_for_queue, retry_cancel_token) - { - let retry_request_id = retry_request.id.clone(); - let retry_attempt = retry_request.retry_attempts; - let retry_key = PriorityHealQueue::make_dedup_key(&retry_request); - let retry_priority = retry_request.priority; - let retry_panic_heal_type = retry_request.heal_type.clone(); - let retry_panic_set_label = heal_request_set_metric_label(&retry_request); - let retry_active_heals = active_heals_clone.clone(); - let retry_heal_queue = heal_queue_clone.clone(); - let retrying_heals_for_spawn = retrying_heals_clone.clone(); - let retry_task_aliases = task_aliases_clone.clone(); - let retry_mrf_repair_notice_targets = mrf_repair_notice_targets_clone.clone(); - let retry_completed_heals = completed_heals_clone.clone(); - let retry_notify = notify_clone.clone(); - let retry_manager_cancel_token = manager_cancel_token.clone(); - let retry_config = config_for_spawn.clone(); - let retry_panic_id = retry_request_id.clone(); - let retry_panic_state = PanicCleanupState { - active_heals: retry_active_heals.clone(), - heal_queue: retry_heal_queue.clone(), - completed_heals: retry_completed_heals.clone(), - task_aliases: retry_task_aliases.clone(), - retrying_heals: retrying_heals_for_spawn.clone(), - mrf_repair_notice_targets: retry_mrf_repair_notice_targets.clone(), - replacement_recovery_anchors: replacement_recovery_anchors_clone.clone(), - statistics: statistics_clone.clone(), - }; - let retry_panic_cancel_token = retry_cancel_token.clone(); - tokio::spawn(async move { - let retry_child = async move { - #[cfg(test)] - panic_if_armed(SchedulerPanicPoint::RetryChild, &retry_request_id); - loop { - tokio::select! { - _ = retry_cancel_token.cancelled() => { - debug!( - target: "rustfs::heal::manager", - event = EVENT_HEAL_QUEUE_ADMISSION, - component = LOG_COMPONENT_HEAL, - subsystem = LOG_SUBSYSTEM_MANAGER, - request_id = %retry_request_id, - priority = ?retry_priority, - retry_attempt, - result = "retry_cancelled", - "Heal retry admission decided" - ); - return; - } - _ = retry_manager_cancel_token.cancelled() => { - retry_cancel_token.cancel(); - retrying_heals_for_spawn.lock().await.remove(&retry_request_id); - return; - } - _ = sleep(retry_delay) => {} + let mut queue = retry_heal_queue.lock().await; + let admission_decision = + Self::admit_request_to_queue(&mut queue, retry_request.clone(), &retry_config, "retry"); + let admission = admission_decision.result; + let should_notify = matches!(admission, HealAdmissionResult::Accepted) + && retry_config.event_driven_scheduler_enable; + // Publish the terminal synchronously while the + // queue transition is protected. The subsequent + // queue -> retrying handoff retains the lock order + // used by operations_snapshot. + let displaced_terminal = admission_decision + .displaced_request + .as_ref() + .map(|request| record_displaced_terminal(&retry_displaced_terminals, request)); + match admission { + HealAdmissionResult::Accepted => { + // Transfer ownership while holding queue -> retrying, + // matching operations_snapshot's lock order. + #[cfg(test)] + pause_retry_ownership_transition(&retry_request_id, true).await; + retrying_heals_for_spawn.lock().await.remove(&retry_request_id); + let displaced_task_id = admission_decision.displaced_task_id().map(ToOwned::to_owned); + drop(queue); + if let (Some(displaced_task_id), Some(displaced_terminal)) = + (displaced_task_id, displaced_terminal) + { + remove_displaced_task_aliases( + &retry_task_aliases, + &retry_displaced_terminals, + &displaced_task_id, + &displaced_terminal, + ) + .await; + remove_mrf_repair_notice_targets( + &retry_mrf_repair_notice_targets, + &displaced_task_id, + ); } { diff --git a/crates/heal/src/heal/manager/tests.rs b/crates/heal/src/heal/manager/tests.rs index 22fc60908..d81fa3530 100644 --- a/crates/heal/src/heal/manager/tests.rs +++ b/crates/heal/src/heal/manager/tests.rs @@ -84,6 +84,7 @@ async fn process_manager_queue_once(manager: &HealManager) { heal_queue: &manager.heal_queue, active_heals: &manager.active_heals, completed_heals: &manager.completed_heals, + displaced_terminals: &manager.displaced_terminals, task_aliases: &manager.task_aliases, retrying_heals: &manager.retrying_heals, mrf_repair_notice_targets: &manager.mrf_repair_notice_targets, @@ -3006,7 +3007,10 @@ async fn test_high_priority_request_displaces_lower_priority_when_queue_full() { HealAdmissionResult::Accepted ); assert_eq!(manager.get_queue_length().await, 1); - assert!(matches!(manager.get_task_status(&low_id).await, Err(Error::TaskNotFound { .. }))); + assert!(matches!( + manager.get_task_status(&low_id).await, + Ok(HealTaskStatus::Failed { error }) if error.contains("reason=displaced") + )); assert_eq!( manager .get_task_status(&high_id) @@ -3016,6 +3020,263 @@ async fn test_high_priority_request_displaces_lower_priority_when_queue_full() { ); } +#[tokio::test] +async fn displaced_task_remains_queryable() { + let manager = HealManager::new( + Arc::new(MockStorage), + Some(HealConfig { + queue_size: 1, + ..HealConfig::default() + }), + ); + let mut displaced = HealRequest::new( + HealType::Bucket { + bucket: "displaced-bucket".to_string(), + }, + HealOptions::default(), + HealPriority::Low, + ); + displaced.id = "displaced-task".to_string(); + let displaced_id = displaced.id.clone(); + manager + .submit_heal_request(displaced) + .await + .expect("displaced request should queue"); + + let successor = HealRequest::new( + HealType::Bucket { + bucket: "successor-bucket".to_string(), + }, + HealOptions::default(), + HealPriority::High, + ); + manager + .submit_heal_request(successor) + .await + .expect("successor should displace low work"); + + let report = manager + .get_task_report(&displaced_id) + .await + .expect("displaced report should remain queryable"); + assert!(matches!(report.status, HealTaskStatus::Failed { ref error } if error.contains("reason=displaced"))); +} + +#[tokio::test] +async fn displaced_archive_failure_keeps_queryable_terminal() { + let manager = HealManager::new(Arc::new(MockStorage), None); + let mut request = HealRequest::new( + HealType::Bucket { + bucket: "archive-failure".to_string(), + }, + HealOptions::default(), + HealPriority::Low, + ); + request.id = "archive-failure-task".to_string(); + let request_id = request.id.clone(); + // The synchronous sidecar is the authoritative fallback when the normal + // completed-task archive has no entry (the failure window that must not + // turn an Accepted ID into NotFound). + record_displaced_terminal(&manager.displaced_terminals, &request); + assert!(manager.completed_heals.lock().await.is_empty()); + assert!(matches!( + manager.get_task_status(&request_id).await, + Ok(HealTaskStatus::Failed { error }) if error.contains("reason=displaced") + )); +} + +#[tokio::test] +async fn scheduler_retry_displacement_keeps_evicted_task_queryable() { + let manager = Arc::new(HealManager::new( + Arc::new(MockStorage), + Some(HealConfig { + queue_size: 1, + event_driven_scheduler_enable: false, + ..HealConfig::default() + }), + )); + let mut retry_request = HealRequest::object("retry-transition".to_string(), "object".to_string(), None); + retry_request.priority = HealPriority::High; + let retry_id = retry_request.id.clone(); + manager + .submit_heal_request(retry_request) + .await + .expect("retry request should queue"); + + // Process exactly one queue cycle so the retry task is spawned without a + // background scheduler consuming the filler request before the retry wakes. + process_manager_queue_once(&manager).await; + tokio::time::timeout(Duration::from_secs(1), async { + loop { + if manager.retrying_heals.lock().await.contains_key(&retry_id) { + break; + } + tokio::task::yield_now().await; + } + }) + .await + .expect("retry request should enter backoff"); + + let filler = HealRequest::new( + HealType::Bucket { + bucket: "retry-displaced-filler".to_string(), + }, + HealOptions::default(), + HealPriority::Low, + ); + let filler_id = filler.id.clone(); + manager + .submit_heal_request(filler) + .await + .expect("filler request should occupy the queue"); + + tokio::time::timeout(Duration::from_secs(5), async { + loop { + if matches!( + manager.get_task_status(&filler_id).await, + Ok(HealTaskStatus::Failed { ref error }) if error.contains("reason=displaced") + ) { + break; + } + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) + .await + .expect("retry admission should displace the filler request"); + assert_eq!(manager.get_queue_length().await, 1); + assert_eq!( + manager.get_task_status(&retry_id).await.expect("retry should be queued"), + HealTaskStatus::Pending + ); +} + +#[tokio::test] +async fn concurrent_displacers_produce_one_terminal_generation() { + let manager = Arc::new(HealManager::new( + Arc::new(MockStorage), + Some(HealConfig { + queue_size: 1, + ..HealConfig::default() + }), + )); + let mut displaced = HealRequest::new( + HealType::Bucket { + bucket: "concurrent-displaced".to_string(), + }, + HealOptions::default(), + HealPriority::Low, + ); + displaced.id = "concurrent-displaced-task".to_string(); + let displaced_id = displaced.id.clone(); + manager + .submit_heal_request(displaced) + .await + .expect("initial request should queue"); + + let first = HealRequest::new( + HealType::Bucket { + bucket: "concurrent-successor-a".to_string(), + }, + HealOptions::default(), + HealPriority::High, + ); + let second = HealRequest::new( + HealType::Bucket { + bucket: "concurrent-successor-b".to_string(), + }, + HealOptions::default(), + HealPriority::High, + ); + let (first_result, second_result) = tokio::join!(manager.submit_heal_request(first), manager.submit_heal_request(second)); + let accepted = [&first_result, &second_result] + .into_iter() + .filter(|result| matches!(result, Ok(HealAdmissionResult::Accepted))) + .count(); + assert_eq!(accepted, 1, "exactly one concurrent displacer should win the full queue"); + assert!( + first_result.is_ok() && second_result.is_ok(), + "the losing request should receive a typed Full result" + ); + let terminals = lock_displaced_terminals(&manager.displaced_terminals); + assert_eq!(terminals.len(), 1); + assert!(terminals.contains_key(&displaced_id)); +} + +#[tokio::test] +async fn successor_chain_is_bounded_and_authorized() { + let manager = HealManager::new( + Arc::new(MockStorage), + Some(HealConfig { + queue_size: 1, + ..HealConfig::default() + }), + ); + let mut original = HealRequest::new( + HealType::Bucket { + bucket: "authorized-original".to_string(), + }, + HealOptions::default(), + HealPriority::Low, + ); + original.id = "authorized-original-task".to_string(); + let original_id = original.id.clone(); + manager.submit_heal_request(original).await.expect("original should queue"); + let mut duplicate = HealRequest::new( + HealType::Bucket { + bucket: "authorized-original".to_string(), + }, + HealOptions::default(), + HealPriority::Low, + ); + duplicate.id = "authorized-duplicate-task".to_string(); + let duplicate_id = duplicate.id.clone(); + manager + .submit_heal_request(duplicate) + .await + .expect("same-target duplicate should merge"); + let successor = HealRequest::new( + HealType::Bucket { + bucket: "authorized-successor".to_string(), + }, + HealOptions::default(), + HealPriority::High, + ); + let successor_id = successor.id.clone(); + manager.submit_heal_request(successor).await.expect("successor should queue"); + assert!(manager.task_aliases.lock().await.is_empty()); + assert!(matches!(manager.get_task_status(&original_id).await, Ok(HealTaskStatus::Failed { .. }))); + assert!(matches!(manager.get_task_status(&duplicate_id).await, Ok(HealTaskStatus::Failed { .. }))); + assert_eq!( + manager + .get_task_status(&successor_id) + .await + .expect("successor should remain queued"), + HealTaskStatus::Pending + ); +} + +#[tokio::test] +async fn displaced_terminal_expires_after_bounded_ttl() { + let manager = HealManager::new(Arc::new(MockStorage), None); + let mut request = HealRequest::new( + HealType::Bucket { + bucket: "expires".to_string(), + }, + HealOptions::default(), + HealPriority::Low, + ); + request.id = "expires-task".to_string(); + let request_id = request.id.clone(); + record_displaced_terminal(&manager.displaced_terminals, &request); + { + let mut terminals = lock_displaced_terminals(&manager.displaced_terminals); + let entry = + Arc::get_mut(terminals.get_mut(&request_id).expect("terminal should be retained")).expect("test owns terminal entry"); + entry.completed_at = SystemTime::now() - KEEP_HEAL_TASK_STATUS_DURATION - Duration::from_secs(1); + } + assert!(matches!(manager.get_task_status(&request_id).await, Err(Error::TaskNotFound { .. }))); +} + #[tokio::test] async fn test_displacing_registered_mrf_task_drops_notice_ownership() { let storage: Arc = Arc::new(MockStorage); diff --git a/crates/protos/src/generated/proto_gen/node_service.rs b/crates/protos/src/generated/proto_gen/node_service.rs index 6fa01bf1e..3885be4a3 100644 --- a/crates/protos/src/generated/proto_gen/node_service.rs +++ b/crates/protos/src/generated/proto_gen/node_service.rs @@ -722,6 +722,10 @@ pub struct DeleteVersionsResponse { pub errors: ::prost::alloc::vec::Vec<::prost::alloc::string::String>, #[prost(message, optional, tag = "3")] pub error: ::core::option::Option, + /// Senders dual-write the legacy strings and typed entries. Receivers prefer typed entries + /// when present and fall back to strings for peers that predate this field. Code zero means success. + #[prost(message, repeated, tag = "4")] + pub item_errors: ::prost::alloc::vec::Vec, } #[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] pub struct ReadMultipleRequest { diff --git a/crates/protos/src/node.proto b/crates/protos/src/node.proto index 290b93b03..70ca93565 100644 --- a/crates/protos/src/node.proto +++ b/crates/protos/src/node.proto @@ -493,6 +493,9 @@ message DeleteVersionsResponse { bool success = 1; repeated string errors = 2; optional Error error = 3; + // Senders dual-write the legacy strings and typed entries. Receivers prefer typed entries + // when present and fall back to strings for peers that predate this field. Code zero means success. + repeated Error item_errors = 4; } message ReadMultipleRequest { diff --git a/rustfs/src/storage/rpc/node_service/disk.rs b/rustfs/src/storage/rpc/node_service/disk.rs index 5d9fbc969..de80a10f0 100644 --- a/rustfs/src/storage/rpc/node_service/disk.rs +++ b/rustfs/src/storage/rpc/node_service/disk.rs @@ -146,6 +146,29 @@ fn encode_file_info_msgpack(value: &FileInfo) -> std::result::Result, Di encode_msgpack_with_capacity(value, "FileInfo", FILE_INFO_MSGPACK_ENCODE_CAPACITY_HINT) } +fn encode_delete_versions_errors(disk_errors: Vec>) -> (Vec, Vec) { + let mut errors = Vec::with_capacity(disk_errors.len()); + let mut item_errors = Vec::with_capacity(disk_errors.len()); + for error in disk_errors { + match error { + Some(error) => { + let code = match &error { + DiskError::Io(source) if source.kind() == std::io::ErrorKind::NotFound => DiskError::FileNotFound.to_u32(), + _ => error.to_u32(), + }; + let error_info = error.to_string(); + errors.push(error_info.clone()); + item_errors.push(Error { code, error_info }); + } + None => { + errors.push(String::new()); + item_errors.push(Error::default()); + } + } + } + (errors, item_errors) +} + fn encode_msgpack_named(value: &T, value_name: &str) -> std::result::Result, DiskError> { let mut serializer = rmp_serde::Serializer::new(Vec::with_capacity(MSGPACK_ENCODE_CAPACITY_HINT)).with_struct_map(); value @@ -552,6 +575,7 @@ impl NodeService { success: false, errors: Vec::new(), error: Some(DiskError::other(format!("decode FileInfoVersions failed: {err}")).into()), + item_errors: Vec::new(), })); } }; @@ -563,30 +587,26 @@ impl NodeService { success: false, errors: Vec::new(), error: Some(DiskError::other(format!("decode DeleteOptions failed: {err}")).into()), + item_errors: Vec::new(), })); } }; - let errors = disk - .delete_versions(&request.volume, versions, opts) - .await - .into_iter() - .map(|error| match error { - Some(e) => e.to_string(), - None => "".to_string(), - }) - .collect(); + let (errors, item_errors) = + encode_delete_versions_errors(disk.delete_versions(&request.volume, versions, opts).await); Ok(Response::new(DeleteVersionsResponse { success: true, errors, error: None, + item_errors, })) } else { Ok(Response::new(DeleteVersionsResponse { success: false, errors: Vec::new(), error: Some(DiskError::other("cannot find disk".to_string()).into()), + item_errors: Vec::new(), })) } } @@ -1612,8 +1632,8 @@ impl NodeService { mod tests { use super::{ compat_response_json, decode_msgpack_or_json, decode_rename_data_request_file_info, - encode_batch_read_version_response_payloads, encode_file_info_msgpack, encode_msgpack, encode_msgpack_named, - encode_read_multiple_response_payloads, encode_rename_data_response_payloads, + encode_batch_read_version_response_payloads, encode_delete_versions_errors, encode_file_info_msgpack, encode_msgpack, + encode_msgpack_named, encode_read_multiple_response_payloads, encode_rename_data_response_payloads, }; use crate::storage::rpc::node_service::make_server; use crate::storage::storage_api::ReadMultipleResp; @@ -1632,6 +1652,18 @@ mod tests { count: u32, } + #[test] + fn delete_versions_response_dual_writes_typed_item_errors() { + let raw_not_found = super::DiskError::Io(std::io::Error::from(std::io::ErrorKind::NotFound)); + let (errors, item_errors) = encode_delete_versions_errors(vec![Some(raw_not_found), None]); + + assert!(errors[0].starts_with("io error ")); + assert!(errors[1].is_empty()); + assert_eq!(item_errors[0].code, super::DiskError::FileNotFound.to_u32()); + assert_eq!(item_errors[0].error_info, errors[0]); + assert_eq!(item_errors[1].code, 0); + } + #[tokio::test] #[serial] async fn handle_read_version_records_attribution_for_missing_disk() {