Skip to main content

ironflow_engine/
schedule.rs

1//! [`CronSchedule`] -- validated cron expression newtype.
2//!
3//! Wraps [`croner::Cron`] to guarantee that any `CronSchedule` value
4//! holds a syntactically valid cron expression. Construction
5//! is fallible; once built the value is safe to pass to
6//! `tokio_cron_scheduler` without further validation.
7//!
8//! A schedule also carries a [`SchedulePolicy`]: what to do with missed
9//! occurrences, with overlapping runs, and the timezone the expression is
10//! evaluated in.
11
12use std::fmt;
13use std::str::FromStr;
14use std::time::Duration;
15
16use chrono_tz::Tz;
17use croner::Cron;
18use ironflow_store::entities::{
19    MAX_CATCHUP_MAX, MAX_CATCHUP_WINDOW_SECS, MIN_CATCHUP_MAX, MIN_CATCHUP_WINDOW_SECS,
20};
21use serde::{Deserialize, Serialize};
22
23pub use ironflow_store::entities::{CatchupPolicy, OverlapPolicy, SchedulePolicy};
24
25/// A validated cron expression.
26///
27/// Internally wraps a [`croner::Cron`], guaranteeing that the expression
28/// has been parsed and validated at construction time.
29///
30/// [`as_str`](CronSchedule::as_str) returns the original expression
31/// as provided by the user, not the normalized form.
32///
33/// The schedule also carries a [`SchedulePolicy`], set with the `with_*`
34/// builder methods. The policy is not part of the serialized form: a
35/// `CronSchedule` serializes to its raw expression, and deserializing one
36/// gives the default policy.
37///
38/// # Examples
39///
40/// ```
41/// use ironflow_engine::schedule::CronSchedule;
42///
43/// let sched = CronSchedule::new("0 0 * * * *").unwrap();
44/// assert_eq!(sched.as_str(), "0 0 * * * *");
45///
46/// let bad = CronSchedule::new("not a cron");
47/// assert!(bad.is_err());
48/// ```
49#[derive(Debug, Clone)]
50pub struct CronSchedule {
51    inner: Cron,
52    raw: String,
53    policy: SchedulePolicy,
54}
55
56impl CronSchedule {
57    /// Parse and validate a cron expression.
58    ///
59    /// Accepts 5-field (standard) or 6-field (with seconds) expressions,
60    /// as supported by [`croner`].
61    ///
62    /// # Errors
63    ///
64    /// Returns an error string if the expression is syntactically invalid.
65    ///
66    /// # Examples
67    ///
68    /// ```
69    /// use ironflow_engine::schedule::CronSchedule;
70    ///
71    /// assert!(CronSchedule::new("0 */5 * * * *").is_ok());
72    /// assert!(CronSchedule::new("garbage").is_err());
73    /// ```
74    pub fn new(expression: &str) -> Result<Self, String> {
75        let inner = Cron::from_str(expression)
76            .map_err(|e| format!("invalid cron expression '{expression}': {e}"))?;
77        Ok(Self {
78            inner,
79            raw: expression.to_string(),
80            policy: SchedulePolicy::default(),
81        })
82    }
83
84    /// Returns the original cron expression string as provided to [`new`](Self::new).
85    pub fn as_str(&self) -> &str {
86        &self.raw
87    }
88
89    /// Evaluate the expression in an IANA timezone instead of UTC.
90    ///
91    /// Occurrences follow the wall clock of the timezone across daylight
92    /// saving changes: `0 9 * * *` in `Europe/Paris` fires at 9:00 Paris
93    /// time in winter and in summer. An occurrence in an hour skipped in
94    /// spring fires once at the end of the gap; an occurrence in an hour
95    /// repeated in autumn fires once, on its first pass.
96    ///
97    /// # Errors
98    ///
99    /// Returns an error string if `tz` is not a known IANA timezone name.
100    ///
101    /// # Examples
102    ///
103    /// ```
104    /// use ironflow_engine::schedule::CronSchedule;
105    ///
106    /// let sched = CronSchedule::new("0 9 * * *")?.with_timezone("Europe/Paris")?;
107    /// assert_eq!(sched.policy().timezone.name(), "Europe/Paris");
108    ///
109    /// assert!(CronSchedule::new("0 9 * * *")?.with_timezone("Mars/Olympus").is_err());
110    /// # Ok::<(), String>(())
111    /// ```
112    pub fn with_timezone(mut self, tz: &str) -> Result<Self, String> {
113        let parsed: Tz = tz
114            .parse()
115            .map_err(|e| format!("invalid timezone '{tz}': {e}"))?;
116        self.policy.timezone = parsed;
117        Ok(self)
118    }
119
120    /// Set what the schedule does with the occurrences it missed while no
121    /// server fired it. Defaults to [`CatchupPolicy::Latest`].
122    ///
123    /// # Examples
124    ///
125    /// ```
126    /// use ironflow_engine::schedule::{CronSchedule, CatchupPolicy};
127    ///
128    /// let sched = CronSchedule::new("0 * * * *")?.with_catchup(CatchupPolicy::All);
129    /// assert_eq!(sched.policy().catchup, CatchupPolicy::All);
130    /// # Ok::<(), String>(())
131    /// ```
132    pub fn with_catchup(mut self, catchup: CatchupPolicy) -> Self {
133        self.policy.catchup = catchup;
134        self
135    }
136
137    /// Set the most runs created to catch up under [`CatchupPolicy::All`].
138    /// The most recent missed occurrences are kept. Defaults to `10`.
139    ///
140    /// # Panics
141    ///
142    /// Panics if `max` is not between `1` and `1000`.
143    ///
144    /// # Examples
145    ///
146    /// ```
147    /// use ironflow_engine::schedule::CronSchedule;
148    ///
149    /// let sched = CronSchedule::new("0 * * * *")?.with_catchup_max(24);
150    /// assert_eq!(sched.policy().catchup_max, 24);
151    /// # Ok::<(), String>(())
152    /// ```
153    pub fn with_catchup_max(mut self, max: u32) -> Self {
154        assert!(
155            (MIN_CATCHUP_MAX..=MAX_CATCHUP_MAX).contains(&max),
156            "catchup_max must be between {MIN_CATCHUP_MAX} and {MAX_CATCHUP_MAX}, got {max}"
157        );
158        self.policy.catchup_max = max;
159        self
160    }
161
162    /// Set how far back a missed occurrence is still caught up. Older ones
163    /// are dropped. Defaults to one day.
164    ///
165    /// # Panics
166    ///
167    /// Panics if `window` is shorter than one minute or longer than 30 days.
168    ///
169    /// # Examples
170    ///
171    /// ```
172    /// use std::time::Duration;
173    ///
174    /// use ironflow_engine::schedule::CronSchedule;
175    ///
176    /// let sched = CronSchedule::new("0 * * * *")?.with_catchup_window(Duration::from_secs(6 * 3600));
177    /// assert_eq!(sched.policy().catchup_window_secs, 21_600);
178    /// # Ok::<(), String>(())
179    /// ```
180    pub fn with_catchup_window(mut self, window: Duration) -> Self {
181        let secs = window.as_secs();
182        assert!(
183            (u64::from(MIN_CATCHUP_WINDOW_SECS)..=u64::from(MAX_CATCHUP_WINDOW_SECS))
184                .contains(&secs),
185            "catchup window must be between {MIN_CATCHUP_WINDOW_SECS} and {MAX_CATCHUP_WINDOW_SECS} seconds, got {secs}"
186        );
187        self.policy.catchup_window_secs =
188            u32::try_from(secs).expect("catchup window bounded by the assert above");
189        self
190    }
191
192    /// Set what the schedule does when an occurrence comes while one of its
193    /// runs is still active. Defaults to [`OverlapPolicy::Allow`].
194    ///
195    /// # Examples
196    ///
197    /// ```
198    /// use ironflow_engine::schedule::{CronSchedule, OverlapPolicy};
199    ///
200    /// let sched = CronSchedule::new("*/5 * * * *")?.with_overlap(OverlapPolicy::Skip);
201    /// assert_eq!(sched.policy().overlap, OverlapPolicy::Skip);
202    /// # Ok::<(), String>(())
203    /// ```
204    pub fn with_overlap(mut self, overlap: OverlapPolicy) -> Self {
205        self.policy.overlap = overlap;
206        self
207    }
208
209    /// Returns the catch-up, overlap and timezone policy of the schedule.
210    ///
211    /// # Examples
212    ///
213    /// ```
214    /// use ironflow_engine::schedule::{CronSchedule, SchedulePolicy};
215    ///
216    /// let sched = CronSchedule::new("0 * * * *")?;
217    /// assert_eq!(sched.policy(), &SchedulePolicy::default());
218    /// # Ok::<(), String>(())
219    /// ```
220    pub fn policy(&self) -> &SchedulePolicy {
221        &self.policy
222    }
223}
224
225impl fmt::Display for CronSchedule {
226    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
227        f.write_str(&self.raw)
228    }
229}
230
231impl PartialEq for CronSchedule {
232    fn eq(&self, other: &Self) -> bool {
233        self.inner == other.inner && self.policy == other.policy
234    }
235}
236
237impl Eq for CronSchedule {}
238
239impl Serialize for CronSchedule {
240    fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
241    where
242        S: serde::Serializer,
243    {
244        serializer.serialize_str(&self.raw)
245    }
246}
247
248impl<'de> Deserialize<'de> for CronSchedule {
249    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
250    where
251        D: serde::Deserializer<'de>,
252    {
253        let s = String::deserialize(deserializer)?;
254        Self::new(&s).map_err(serde::de::Error::custom)
255    }
256}
257
258#[cfg(test)]
259mod tests {
260    use super::*;
261
262    #[test]
263    fn valid_six_field_expression() {
264        let sched = CronSchedule::new("0 0 * * * *").unwrap();
265        assert_eq!(sched.as_str(), "0 0 * * * *");
266    }
267
268    #[test]
269    fn valid_five_field_expression() {
270        let sched = CronSchedule::new("*/5 * * * *").unwrap();
271        assert_eq!(sched.as_str(), "*/5 * * * *");
272    }
273
274    #[test]
275    fn valid_complex_expression() {
276        let sched = CronSchedule::new("0 30 9 * * MON-FRI").unwrap();
277        assert_eq!(sched.as_str(), "0 30 9 * * MON-FRI");
278    }
279
280    #[test]
281    fn invalid_expression_returns_error() {
282        let result = CronSchedule::new("not a cron");
283        assert!(result.is_err());
284        let err = result.unwrap_err();
285        assert!(err.contains("invalid cron expression"));
286    }
287
288    #[test]
289    fn empty_expression_returns_error() {
290        assert!(CronSchedule::new("").is_err());
291    }
292
293    #[test]
294    fn display_shows_original_expression() {
295        let sched = CronSchedule::new("0 0 12 * * *").unwrap();
296        assert_eq!(format!("{sched}"), "0 0 12 * * *");
297    }
298
299    #[test]
300    fn semantic_equality() {
301        let a = CronSchedule::new("0 0 * * * MON-FRI").unwrap();
302        let b = CronSchedule::new("0 0 * * * 1-5").unwrap();
303        assert_eq!(a, b);
304    }
305
306    #[test]
307    fn inequality_on_different_expressions() {
308        let a = CronSchedule::new("0 0 * * * *").unwrap();
309        let b = CronSchedule::new("0 30 * * * *").unwrap();
310        assert_ne!(a, b);
311    }
312
313    #[test]
314    fn serde_roundtrip() {
315        let sched = CronSchedule::new("0 */5 * * * *").unwrap();
316        let json = serde_json::to_string(&sched).unwrap();
317        assert_eq!(json, "\"0 */5 * * * *\"");
318        let back: CronSchedule = serde_json::from_str(&json).unwrap();
319        assert_eq!(back, sched);
320    }
321
322    #[test]
323    fn deserialize_invalid_expression_fails() {
324        let result: Result<CronSchedule, _> = serde_json::from_str("\"garbage\"");
325        assert!(result.is_err());
326    }
327
328    #[test]
329    fn default_policy_is_latest_allow_utc() {
330        let sched = CronSchedule::new("0 * * * *").unwrap();
331        assert_eq!(sched.policy().catchup, CatchupPolicy::Latest);
332        assert_eq!(sched.policy().overlap, OverlapPolicy::Allow);
333        assert_eq!(sched.policy().timezone.name(), "UTC");
334        assert_eq!(sched.policy(), &SchedulePolicy::default());
335    }
336
337    #[test]
338    fn with_timezone_accepts_iana_name() {
339        let sched = CronSchedule::new("0 9 * * *")
340            .unwrap()
341            .with_timezone("Europe/Paris")
342            .unwrap();
343        assert_eq!(sched.policy().timezone.name(), "Europe/Paris");
344        assert_eq!(sched.as_str(), "0 9 * * *");
345    }
346
347    #[test]
348    fn with_timezone_rejects_unknown_name() {
349        let err = CronSchedule::new("0 9 * * *")
350            .unwrap()
351            .with_timezone("Mars/Olympus")
352            .unwrap_err();
353        assert!(err.contains("invalid timezone 'Mars/Olympus'"), "{err}");
354    }
355
356    #[test]
357    fn builders_set_the_policy() {
358        let sched = CronSchedule::new("0 * * * *")
359            .unwrap()
360            .with_catchup(CatchupPolicy::All)
361            .with_catchup_max(1000)
362            .with_catchup_window(Duration::from_secs(60))
363            .with_overlap(OverlapPolicy::Skip);
364        let policy = sched.policy();
365        assert_eq!(policy.catchup, CatchupPolicy::All);
366        assert_eq!(policy.catchup_max, 1000);
367        assert_eq!(policy.catchup_window_secs, 60);
368        assert_eq!(policy.overlap, OverlapPolicy::Skip);
369    }
370
371    #[test]
372    #[should_panic(expected = "catchup_max must be between 1 and 1000")]
373    fn with_catchup_max_zero_panics() {
374        let _ = CronSchedule::new("0 * * * *").unwrap().with_catchup_max(0);
375    }
376
377    #[test]
378    #[should_panic(expected = "catchup window must be between 60 and 2592000 seconds")]
379    fn with_catchup_window_below_a_minute_panics() {
380        let _ = CronSchedule::new("0 * * * *")
381            .unwrap()
382            .with_catchup_window(Duration::from_secs(59));
383    }
384
385    #[test]
386    #[should_panic(expected = "catchup window must be between 60 and 2592000 seconds")]
387    fn with_catchup_window_above_thirty_days_panics() {
388        let _ = CronSchedule::new("0 * * * *")
389            .unwrap()
390            .with_catchup_window(Duration::from_secs(2_592_001));
391    }
392
393    #[test]
394    fn policies_take_part_in_equality() {
395        let a = CronSchedule::new("0 * * * *").unwrap();
396        let b = CronSchedule::new("0 * * * *")
397            .unwrap()
398            .with_overlap(OverlapPolicy::Skip);
399        let c = CronSchedule::new("0 * * * *")
400            .unwrap()
401            .with_timezone("Europe/Paris")
402            .unwrap();
403        assert_ne!(a, b);
404        assert_ne!(a, c);
405        assert_eq!(b, b.clone());
406    }
407
408    #[test]
409    fn serialized_form_drops_the_policy() {
410        let sched = CronSchedule::new("0 9 * * *")
411            .unwrap()
412            .with_catchup(CatchupPolicy::Skip);
413        let json = serde_json::to_string(&sched).unwrap();
414        assert_eq!(json, "\"0 9 * * *\"");
415        let back: CronSchedule = serde_json::from_str(&json).unwrap();
416        assert_eq!(back.policy(), &SchedulePolicy::default());
417    }
418}