1use crate::contracts::{
2 JobContract, JobEvent, JobEventListener, JobState, JobStore, JobType, MetricsExporter,
3 RetryPolicy,
4};
5use crate::{CronError, CronResult};
6use chrono::{DateTime, Utc};
7use std::borrow::Cow;
8use std::sync::Arc;
9use tokio::time::{sleep, timeout};
10
11#[derive(Clone)]
17pub struct JobItem {
18 job: Arc<dyn JobContract>,
19 listeners: Vec<Arc<dyn JobEventListener>>,
20 metrics_exporter: Option<Arc<dyn MetricsExporter>>,
21 job_store: Option<Arc<dyn JobStore>>,
22}
23
24impl std::fmt::Debug for JobItem {
25 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
26 f.debug_struct("JobItem")
27 .field("id", &self.id())
28 .field("name", &self.name())
29 .field("job_type", &self.job_type())
30 .field("listeners_len", &self.listeners.len())
31 .field("metrics_exporter", &self.metrics_exporter.is_some())
32 .field("job_store", &self.job_store.is_some())
33 .finish()
34 }
35}
36
37impl JobItem {
38 pub fn new(
43 job: Arc<dyn JobContract>,
44 listeners: Vec<Arc<dyn JobEventListener>>,
45 metrics_exporter: Option<Arc<dyn MetricsExporter>>,
46 job_store: Option<Arc<dyn JobStore>>,
47 ) -> CronResult<Self> {
48 let _ = job.schedule();
50 Ok(JobItem {
51 job,
52 listeners,
53 metrics_exporter,
54 job_store,
55 })
56 }
57
58 #[allow(dead_code)]
60 pub fn id(&self) -> Cow<'_, str> {
61 self.job.id()
62 }
63
64 pub fn name(&self) -> Cow<'_, str> {
66 self.job.name()
67 }
68
69 pub fn job_type(&self) -> JobType {
71 self.job.job_type()
72 }
73
74 pub fn next_run_time(&self) -> Option<DateTime<Utc>> {
76 match self.job.job_type() {
77 JobType::Once => {
78 let run_at = self.job.run_at()?;
79 if run_at > Utc::now() {
80 Some(run_at)
81 } else {
82 None
83 }
84 }
85 JobType::Recurring => {
86 let mut after = Utc::now();
87 if let Some(start_after) = self.job.start_after()
88 && start_after > after
89 {
90 after = start_after;
91 }
92 self.job.schedule().next_after(&after, self.job.timezone())
93 }
94 }
95 }
96
97 pub fn next_run_after(&self, after: DateTime<Utc>) -> Option<DateTime<Utc>> {
99 match self.job.job_type() {
100 JobType::Once => {
101 let run_at = self.job.run_at()?;
102 if run_at > after { Some(run_at) } else { None }
103 }
104 JobType::Recurring => {
105 let mut after = after;
106 if let Some(start_after) = self.job.start_after()
107 && start_after > after
108 {
109 after = start_after;
110 }
111 self.job.schedule().next_after(&after, self.job.timezone())
112 }
113 }
114 }
115
116 pub fn priority(&self) -> i32 {
118 self.job.priority()
119 }
120
121 pub fn concurrency_limit(&self) -> Option<usize> {
123 self.job.concurrency_limit()
124 }
125
126 pub fn misfire_policy(&self) -> crate::contracts::MisfirePolicy {
128 self.job.misfire_policy()
129 }
130
131 async fn emit_event(&self, event: JobEvent) {
132 for listener in &self.listeners {
133 listener.on_event(event.clone()).await;
134 }
135 }
136
137 pub async fn run(&self) -> CronResult<()> {
140 let retry_policy = self.job.retry_policy();
141 let mut attempts = 0;
142 let id = self.id().to_string();
143 let name = self.name().to_string();
144
145 let mut state = if let Some(store) = &self.job_store {
146 store.get_state(&id).await?.unwrap_or_default()
147 } else {
148 JobState::default()
149 };
150
151 loop {
152 self.emit_event(JobEvent::Started {
153 id: id.clone(),
154 name: name.clone(),
155 })
156 .await;
157
158 if let Some(exporter) = &self.metrics_exporter {
159 exporter.record_start(&id, &name);
160 }
161
162 self.job.on_start().await;
163 let start_time = Utc::now();
164 state.last_run = Some(start_time);
165
166 let result = if let Some(duration) = self.job.timeout() {
167 match timeout(duration, self.job.run()).await {
168 Ok(res) => res,
169 Err(_) => Err(CronError::ExecutionError(anyhow::anyhow!(
170 "Job timed out after {:?}",
171 duration
172 ))),
173 }
174 } else {
175 self.job.run().await
176 };
177
178 match result {
179 Ok(()) => {
180 let end_time = Utc::now();
181 let duration = (end_time - start_time).to_std().unwrap_or_default();
182
183 state.last_success = Some(end_time);
184 state.consecutive_failures = 0;
185 if let Some(store) = &self.job_store {
186 store.save_state(&id, &state).await?;
187 }
188
189 self.emit_event(JobEvent::Completed {
190 id: id.clone(),
191 name: name.clone(),
192 duration,
193 })
194 .await;
195
196 if let Some(exporter) = &self.metrics_exporter {
197 exporter.record_completion(&id, &name, duration);
198 }
199
200 self.job.on_complete().await;
201 return Ok(());
202 }
203 Err(err) => {
204 attempts += 1;
205 state.last_failure = Some(Utc::now());
206 state.consecutive_failures += 1;
207 if let Some(store) = &self.job_store {
208 store.save_state(&id, &state).await?;
209 }
210
211 let should_retry = match &retry_policy {
212 RetryPolicy::None => false,
213 RetryPolicy::Fixed { max_retries, .. } => attempts <= *max_retries,
214 RetryPolicy::Exponential { max_retries, .. } => attempts <= *max_retries,
215 };
216
217 if should_retry {
218 let delay = match &retry_policy {
219 RetryPolicy::None => std::time::Duration::from_secs(0),
220 RetryPolicy::Fixed { interval, .. } => *interval,
221 RetryPolicy::Exponential {
222 initial_interval,
223 max_interval,
224 ..
225 } => {
226 let exponent = ((attempts as i32 - 1) as u32).min(1023);
228 let backoff =
229 initial_interval.as_secs_f64() * (2.0f64.powi(exponent as i32));
230 std::time::Duration::from_secs_f64(backoff).min(*max_interval)
231 }
232 };
233
234 self.emit_event(JobEvent::Retrying {
235 id: id.clone(),
236 name: name.clone(),
237 attempt: attempts,
238 delay,
239 })
240 .await;
241
242 if let Some(exporter) = &self.metrics_exporter {
243 exporter.record_retry(&id, &name);
244 }
245
246 tracing::warn!(
247 "[{}] Job failed, retrying in {:?} (attempt {}): {:?}",
248 self.name(),
249 delay,
250 attempts,
251 err
252 );
253
254 sleep(delay).await;
255 continue;
256 } else {
257 self.emit_event(JobEvent::Failed {
258 id: id.clone(),
259 name: name.clone(),
260 error: err.to_string(),
261 })
262 .await;
263
264 if let Some(exporter) = &self.metrics_exporter {
265 exporter.record_failure(&id, &name);
266 }
267
268 self.job.on_error(&err).await;
269 return Err(err);
270 }
271 }
272 }
273 }
274 }
275}