Skip to main content

chronon_runtime/
retry.rs

1//! Finalize failed/timed-out runs and enqueue delayed retries.
2//!
3//! # Concern → API
4//!
5//! | Concern | API |
6//! |---------|-----|
7//! | Persist failure + schedule retry | [`finalize_failed_run`] |
8//! | Policy decode / backoff | [`chronon_core::RetryPolicy`] on the job |
9//!
10//! # Examples
11//!
12//! ```ignore
13//! use chronon_core::models::{RetryPolicy, RunStatus};
14//! use chronon_runtime::retry::finalize_failed_run;
15//!
16//! // After a worker script returns Err or times out:
17//! finalize_failed_run(
18//!     &store,
19//!     run,
20//!     &job,
21//!     RunStatus::Failed,
22//!     "script exploded",
23//!     None,
24//! ).await;
25//! // When `job.retry_policy().should_retry(run.attempt)`, a new Queued run
26//! // is created with `attempt + 1` and `scheduled_for = now + delay`.
27//! ```
28
29use std::sync::Arc;
30
31use chrono::{Duration, Utc};
32use chronon_core::models::{Job, Run, RunStatus};
33use chronon_core::store::SchedulerStore;
34use chronon_core::{redact_credentials_in_text, sanitize_error_message};
35use chronon_telemetry::CapturedLogs;
36use tracing::{info, warn};
37
38/// Persist a terminal failure/timeout and enqueue a retry when the job policy allows.
39///
40/// Applies optional [`CapturedLogs`] to the failed run before update. Does **not** retry
41/// `Canceled` or `Success` (callers must only pass [`RunStatus::Failed`] or
42/// [`RunStatus::Timeout`]).
43pub async fn finalize_failed_run(
44    store: &Arc<dyn SchedulerStore>,
45    mut run: Run,
46    job: &Job,
47    status: RunStatus,
48    error: impl Into<String>,
49    logs: Option<CapturedLogs>,
50) {
51    let message = sanitize_error_message(&redact_credentials_in_text(&error.into()));
52    match status {
53        RunStatus::Timeout => run.timeout(message.clone()),
54        _ => run.fail(message.clone()),
55    }
56    if let Some(mut captured) = logs {
57        captured.ensure_stderr_message(&message);
58        run.stdout_text = captured.stdout_text;
59        run.stderr_text = captured.stderr_text;
60    } else {
61        run.stderr_text = Some(message.clone());
62    }
63
64    if let Err(e) = store.update_run(&run).await {
65        warn!(
66            run_id = %run.run_id,
67            job_id = ?run.job_id,
68            error = %e,
69            "failed to persist terminal run before retry decision"
70        );
71        return;
72    }
73
74    let policy = job.retry_policy();
75    if !policy.should_retry(run.attempt) {
76        info!(
77            run_id = %run.run_id,
78            attempt = run.attempt,
79            max_attempts = policy.max_attempts,
80            "retry policy exhausted or disabled; no further attempt"
81        );
82        return;
83    }
84
85    let delay_ms = policy.delay_ms_after(run.attempt) as i64;
86    let scheduled_for = Utc::now() + Duration::milliseconds(delay_ms.max(0));
87    let mut next = Run::for_job(
88        run.job_id.clone().unwrap_or_default(),
89        &run.script_name,
90        scheduled_for,
91    );
92    next.attempt = run.attempt + 1;
93    next.actor_json = run.actor_json.clone();
94    next.params_json = run.params_json.clone();
95    next.pool_id = run.pool_id.clone();
96    next.placement_json = run.placement_json.clone();
97    next.parent_run_id = run.parent_run_id.clone();
98    next.root_run_id = run.root_run_id.clone();
99
100    match store.create_run(&next).await {
101        Ok(()) => {
102            info!(
103                prior_run_id = %run.run_id,
104                next_run_id = %next.run_id,
105                next_attempt = next.attempt,
106                delay_ms,
107                "enqueued retry run"
108            );
109        }
110        Err(e) => {
111            warn!(
112                prior_run_id = %run.run_id,
113                next_attempt = next.attempt,
114                error = %e,
115                "failed to create retry run"
116            );
117        }
118    }
119}
120
121#[cfg(test)]
122mod tests {
123    #![allow(clippy::unwrap_used, clippy::expect_used)]
124
125    use super::*;
126    use chronon_backend_mem::InMemorySchedulerStore;
127    use chronon_core::models::{Job, RetryPolicy, RunStatus};
128    use chronon_core::store::SchedulerStore;
129
130    fn job_with_retries(max_attempts: u32, base_delay_ms: u64) -> Job {
131        let mut job = Job::new("retry-job", "script");
132        job.set_retry_policy(&RetryPolicy {
133            max_attempts,
134            base_delay_ms,
135            backoff_multiplier: 2.0,
136            max_delay_ms: 10_000,
137        });
138        job
139    }
140
141    #[tokio::test]
142    async fn failed_run_schedules_retry_with_incremented_attempt() {
143        let store: Arc<dyn SchedulerStore> = Arc::new(InMemorySchedulerStore::new());
144        let job = job_with_retries(3, 100);
145        store.upsert_job(&job).await.unwrap();
146
147        let run = Run::for_job(&job.job_id, "script", Utc::now());
148        assert_eq!(run.attempt, 1);
149        store.create_run(&run).await.unwrap();
150
151        finalize_failed_run(
152            &store,
153            run.clone(),
154            &job,
155            RunStatus::Failed,
156            "boom",
157            Some(CapturedLogs {
158                stdout_text: Some("out".into()),
159                stderr_text: None,
160            }),
161        )
162        .await;
163
164        let failed = store.get_run(&run.run_id).await.unwrap().unwrap();
165        assert_eq!(failed.status, RunStatus::Failed);
166        assert_eq!(failed.stdout_text.as_deref(), Some("out"));
167        assert_eq!(failed.stderr_text.as_deref(), Some("boom"));
168
169        let runs = store.list_runs_for_job(&job.job_id, 10).await.unwrap();
170        assert_eq!(runs.len(), 2);
171        let next = runs
172            .iter()
173            .find(|r| r.run_id != run.run_id)
174            .expect("retry run");
175        assert_eq!(next.attempt, 2);
176        assert_eq!(next.status, RunStatus::Queued);
177        let delta = (next.scheduled_for - Utc::now()).num_milliseconds();
178        assert!((50..=250).contains(&delta), "delay ~100ms, got {delta}");
179    }
180
181    #[tokio::test]
182    async fn timeout_also_retries() {
183        let store: Arc<dyn SchedulerStore> = Arc::new(InMemorySchedulerStore::new());
184        let job = job_with_retries(1, 0);
185        store.upsert_job(&job).await.unwrap();
186        let run = Run::for_job(&job.job_id, "script", Utc::now());
187        store.create_run(&run).await.unwrap();
188
189        finalize_failed_run(&store, run, &job, RunStatus::Timeout, "timeout", None).await;
190
191        let runs = store.list_runs_for_job(&job.job_id, 10).await.unwrap();
192        assert_eq!(runs.len(), 2);
193        assert!(runs.iter().any(|r| r.status == RunStatus::Timeout));
194        assert!(runs
195            .iter()
196            .any(|r| r.status == RunStatus::Queued && r.attempt == 2));
197    }
198
199    #[tokio::test]
200    async fn exhausted_attempts_do_not_enqueue() {
201        let store: Arc<dyn SchedulerStore> = Arc::new(InMemorySchedulerStore::new());
202        let job = job_with_retries(1, 0);
203        store.upsert_job(&job).await.unwrap();
204        let mut run = Run::for_job(&job.job_id, "script", Utc::now());
205        run.attempt = 2; // first retry already used (max_attempts=1)
206        store.create_run(&run).await.unwrap();
207
208        finalize_failed_run(&store, run, &job, RunStatus::Failed, "done", None).await;
209
210        let runs = store.list_runs_for_job(&job.job_id, 10).await.unwrap();
211        assert_eq!(runs.len(), 1);
212        assert_eq!(runs[0].status, RunStatus::Failed);
213    }
214
215    #[tokio::test]
216    async fn sanitize_strips_secret_prefixes_before_persist() {
217        let store: Arc<dyn SchedulerStore> = Arc::new(InMemorySchedulerStore::new());
218        let job = Job::new("sanitize-job", "script");
219        store.upsert_job(&job).await.unwrap();
220        let run = Run::for_job(&job.job_id, "script", Utc::now());
221        store.create_run(&run).await.unwrap();
222
223        finalize_failed_run(
224            &store,
225            run.clone(),
226            &job,
227            RunStatus::Failed,
228            "handler failed password=hunter2 leftover",
229            None,
230        )
231        .await;
232
233        let failed = store.get_run(&run.run_id).await.unwrap().unwrap();
234        let stderr = failed.stderr_text.as_deref().unwrap_or("");
235        assert!(stderr.contains("[redacted]"));
236        assert!(!stderr.contains("hunter2"));
237    }
238
239    #[tokio::test]
240    async fn default_policy_does_not_retry() {
241        let store: Arc<dyn SchedulerStore> = Arc::new(InMemorySchedulerStore::new());
242        let job = Job::new("no-retry", "script");
243        store.upsert_job(&job).await.unwrap();
244        let run = Run::for_job(&job.job_id, "script", Utc::now());
245        store.create_run(&run).await.unwrap();
246
247        finalize_failed_run(&store, run, &job, RunStatus::Failed, "x", None).await;
248
249        assert_eq!(
250            store
251                .list_runs_for_job(&job.job_id, 10)
252                .await
253                .unwrap()
254                .len(),
255            1
256        );
257    }
258}