Skip to main content

foxtive_cron/
job.rs

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/// An internal wrapper around a `JobContract` that caches the parsed schedule
12/// and exposes helper methods used by the scheduler.
13///
14/// Constructed once per registered job via [`JobItem::new`], ensuring the
15/// schedule is valid at registration time rather than at execution time.
16#[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    /// Wrap a [`JobContract`] implementor.
39    ///
40    /// The schedule is validated via [`JobContract::schedule`] at this point.
41    /// Returns an error if the job's schedule is invalid.
42    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        // Eagerly access the schedule to trigger any validation at registration time.
49        let _ = job.schedule();
50        Ok(JobItem {
51            job,
52            listeners,
53            metrics_exporter,
54            job_store,
55        })
56    }
57
58    /// The job's stable unique identifier.
59    #[allow(dead_code)]
60    pub fn id(&self) -> Cow<'_, str> {
61        self.job.id()
62    }
63
64    /// The job's human-readable name.
65    pub fn name(&self) -> Cow<'_, str> {
66        self.job.name()
67    }
68
69    /// The type of job.
70    pub fn job_type(&self) -> JobType {
71        self.job.job_type()
72    }
73
74    /// Computes the next scheduled execution time from now.
75    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    /// Computes the next scheduled execution time after a specific time.
98    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    /// Returns the job's priority.
117    pub fn priority(&self) -> i32 {
118        self.job.priority()
119    }
120
121    /// Returns the job's concurrency limit.
122    pub fn concurrency_limit(&self) -> Option<usize> {
123        self.job.concurrency_limit()
124    }
125
126    /// Returns the job's misfire policy.
127    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    /// Runs the lifecycle sequence: `on_start` → `run` → `on_complete` / `on_error`.
138    /// Handles timeouts and retries internally.
139    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                                // Cap the exponent to prevent f64 overflow (2^1023 is near f64::MAX)
227                                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}