1use 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
38pub 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; 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}