ironflow_engine/
schedule.rs1use 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#[derive(Debug, Clone)]
50pub struct CronSchedule {
51 inner: Cron,
52 raw: String,
53 policy: SchedulePolicy,
54}
55
56impl CronSchedule {
57 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 pub fn as_str(&self) -> &str {
86 &self.raw
87 }
88
89 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 pub fn with_catchup(mut self, catchup: CatchupPolicy) -> Self {
133 self.policy.catchup = catchup;
134 self
135 }
136
137 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 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 pub fn with_overlap(mut self, overlap: OverlapPolicy) -> Self {
205 self.policy.overlap = overlap;
206 self
207 }
208
209 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}