diff --git a/crates/heal/src/error.rs b/crates/heal/src/error.rs index 68ca7eef4..413f02c54 100644 --- a/crates/heal/src/error.rs +++ b/crates/heal/src/error.rs @@ -82,8 +82,8 @@ impl Error { /// Whether a heal operation can be retried without changing its inputs. pub(crate) fn is_recoverable_heal(&self) -> bool { match self { - Error::TaskCancelled => false, - Error::TaskTimeout | Error::TransientSkip { .. } => true, + Error::TaskCancelled | Error::TaskTimeout => false, + Error::TransientSkip { .. } => true, Error::Storage(err) => { err.is_quorum_error() || matches!( @@ -165,4 +165,9 @@ mod tests { assert!(Error::Storage(EcstoreError::DiskNotFound).is_recoverable_heal()); assert!(Error::Storage(EcstoreError::VolumeNotFound).is_recoverable_heal()); } + + #[test] + fn task_timeout_is_terminal() { + assert!(!Error::TaskTimeout.is_recoverable_heal()); + } } diff --git a/crates/heal/src/heal/manager.rs b/crates/heal/src/heal/manager.rs index fa831b129..79ba6b169 100644 --- a/crates/heal/src/heal/manager.rs +++ b/crates/heal/src/heal/manager.rs @@ -673,6 +673,12 @@ fn retry_request_for_result(task: &HealTask, result: &Result<()>) -> Option<(Hea Some((request, delay, error)) } +async fn retry_request_for_result_with_budget(task: &HealTask, result: &Result<()>) -> Option<(HealRequest, Duration, String)> { + let (_, delay, error) = retry_request_for_result(task, result)?; + let request = task.retry_request_with_remaining_timeout().await.ok()?; + Some((request, delay, error)) +} + fn recoverable_heal_retry_delay(retry_attempt: u32) -> Duration { let retry_attempt = retry_attempt.clamp(1, 5); let delay = Duration::from_secs(2_u64.saturating_pow(retry_attempt)); @@ -690,7 +696,7 @@ pub struct HealConfig { pub max_concurrent_heals: usize, /// Maximum concurrent heal tasks allowed for a single erasure set pub max_concurrent_per_set: usize, - /// Task timeout + /// Aggregate task execution timeout across recoverable retries pub task_timeout: Duration, /// Queue size pub queue_size: usize, @@ -3106,7 +3112,7 @@ impl HealManager { "Heal scheduler task started" ); let result = task.execute().await; - let retry_request = retry_request_for_result(task.as_ref(), &result); + let retry_request = retry_request_for_result_with_budget(task.as_ref(), &result).await; match &result { Ok(_) => { debug!( @@ -4539,6 +4545,25 @@ mod tests { assert!(retry_error.contains("Lock acquisition timeout")); } + #[tokio::test] + async fn retry_request_for_result_preserves_remaining_timeout_budget() { + let storage: Arc = Arc::new(MockStorage); + let mut request = HealRequest::object("retry-transition".to_string(), "object".to_string(), None); + request.options.timeout = Some(Duration::from_secs(60)); + let task = HealTask::from_request(request, storage); + let result = task.execute().await; + + let (retry_request, _, _) = retry_request_for_result_with_budget(&task, &result) + .await + .expect("read quorum failure should retain the unused timeout budget"); + let remaining = retry_request + .options + .timeout + .expect("configured timeout should remain present"); + assert!(remaining < Duration::from_secs(60)); + assert!(remaining > Duration::from_secs(59)); + } + #[test] fn test_retry_request_for_incomplete_heal_rename() { let storage: Arc = Arc::new(MockStorage); @@ -6054,7 +6079,7 @@ mod tests { process_manager_queue_once(&manager).await; let defaulted_status = tokio::time::timeout(Duration::from_secs(1), async { loop { - if let Ok(status @ HealTaskStatus::Retrying { .. }) = manager.get_task_status(&defaulted_id).await { + if let Ok(status @ HealTaskStatus::Timeout) = manager.get_task_status(&defaulted_id).await { break status; } tokio::task::yield_now().await; @@ -6062,23 +6087,8 @@ mod tests { }) .await .expect("configured timeout should finish the task"); - assert!(matches!(defaulted_status, HealTaskStatus::Retrying { .. })); - assert_eq!( - manager - .retrying_heals - .lock() - .await - .get(&defaulted_id) - .expect("timed out task should retain its retry request") - .request - .options - .timeout, - Some(Duration::ZERO) - ); - manager - .cancel_task(&defaulted_id) - .await - .expect("retrying timeout task should be cancelled"); + assert_eq!(defaulted_status, HealTaskStatus::Timeout); + assert!(manager.retrying_heals.lock().await.get(&defaulted_id).is_none()); let mut explicit = bucket_request("explicit-timeout", HealPriority::Normal, HealRequestSource::Admin); explicit.options.timeout = Some(Duration::from_secs(60)); diff --git a/crates/heal/src/heal/task.rs b/crates/heal/src/heal/task.rs index 916b51822..62123c418 100644 --- a/crates/heal/src/heal/task.rs +++ b/crates/heal/src/heal/task.rs @@ -196,7 +196,7 @@ pub struct HealOptions { /// Whether to skip namespace locking #[serde(default)] pub no_lock: bool, - /// Timeout + /// Aggregate execution timeout across recoverable manager retries pub timeout: Option, /// pool index pub pool_index: Option, @@ -442,6 +442,14 @@ impl HealTask { } } + pub(crate) async fn retry_request_with_remaining_timeout(&self) -> Result { + let mut request = self.retry_request(); + if self.options.timeout.is_some() { + request.options.timeout = self.remaining_timeout().await?; + } + Ok(request) + } + pub(crate) fn from_replacement_recovery_request( request: HealRequest, storage: Arc, @@ -2657,6 +2665,36 @@ mod tests { use super::super::storage_api::status::BucketInfo; + #[tokio::test] + async fn retry_request_carries_remaining_timeout_budget() { + let storage: Arc = Arc::new(MockStorage::default()); + let mut request = HealRequest::bucket("bucket".to_string()); + request.options.timeout = Some(Duration::from_secs(100)); + let task = HealTask::from_request(request, storage.clone()); + *task.task_start_instant.write().await = Some(Instant::now() - Duration::from_secs(40)); + + let retry = task + .retry_request_with_remaining_timeout() + .await + .expect("first retry should retain the unused timeout budget"); + let first_remaining = retry.options.timeout.expect("configured timeout should remain present"); + assert!(first_remaining <= Duration::from_secs(60)); + assert!(first_remaining > Duration::from_secs(59)); + + let retry_task = HealTask::from_request(retry, storage); + *retry_task.task_start_instant.write().await = Some(Instant::now() - Duration::from_secs(20)); + let second_retry = retry_task + .retry_request_with_remaining_timeout() + .await + .expect("second retry should retain only the unused aggregate budget"); + let second_remaining = second_retry + .options + .timeout + .expect("configured timeout should remain present"); + assert!(second_remaining <= Duration::from_secs(40)); + assert!(second_remaining > Duration::from_secs(39)); + } + #[test] fn format_result_requires_every_requested_target_to_be_ok() { let result = HealResultItem {