ironflow_store/entities/schedule.rs
1//! Schedule entity for periodic workflow execution.
2
3use std::collections::HashMap;
4
5use chrono::{DateTime, SecondsFormat, Utc};
6use chrono_tz::Tz;
7use serde::{Deserialize, Serialize};
8use serde_json::Value;
9use strum::{AsRefStr, Display, EnumString, IntoStaticStr};
10use uuid::Uuid;
11
12use super::{NewRun, RunActor, RunCreation, TriggerKind};
13
14/// Where a schedule was created.
15///
16/// `Handler` schedules are declared in code via [`WorkflowHandler::schedule()`]
17/// and synced to the database at startup. `Api` schedules are created by users
18/// through the REST API, CLI, or dashboard.
19///
20/// # Examples
21///
22/// ```
23/// use ironflow_store::entities::ScheduleSource;
24///
25/// let source = ScheduleSource::Handler;
26/// assert_eq!(source.as_str(), "handler");
27///
28/// let parsed: ScheduleSource = "api".parse().unwrap();
29/// assert_eq!(parsed, ScheduleSource::Api);
30/// ```
31#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
32#[derive(
33 Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, Display, EnumString, IntoStaticStr,
34)]
35#[serde(rename_all = "snake_case")]
36#[strum(serialize_all = "snake_case")]
37pub enum ScheduleSource {
38 /// Declared in code via `WorkflowHandler::schedule()`.
39 Handler,
40 /// Created via the REST API.
41 Api,
42}
43
44impl ScheduleSource {
45 /// String representation used in the database.
46 pub fn as_str(&self) -> &'static str {
47 self.into()
48 }
49}
50
51/// What a schedule does with the occurrences it missed while no server was
52/// firing it (downtime, a long deploy, a stalled ticker).
53///
54/// # Examples
55///
56/// ```
57/// use ironflow_store::entities::CatchupPolicy;
58///
59/// assert_eq!(CatchupPolicy::default(), CatchupPolicy::Latest);
60/// assert_eq!(CatchupPolicy::All.as_ref(), "all");
61///
62/// let parsed: CatchupPolicy = "skip".parse().unwrap();
63/// assert_eq!(parsed, CatchupPolicy::Skip);
64/// ```
65#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
66#[derive(
67 Debug,
68 Clone,
69 Copy,
70 Default,
71 PartialEq,
72 Eq,
73 Serialize,
74 Deserialize,
75 Display,
76 EnumString,
77 IntoStaticStr,
78 AsRefStr,
79)]
80#[serde(rename_all = "snake_case")]
81#[strum(serialize_all = "snake_case")]
82pub enum CatchupPolicy {
83 /// Run only the most recent missed occurrence.
84 #[default]
85 Latest,
86 /// Run every missed occurrence, oldest first, up to
87 /// [`SchedulePolicy::catchup_max`].
88 All,
89 /// Run no late occurrence: only an occurrence fired on time runs.
90 Skip,
91}
92
93/// What a schedule does when an occurrence comes while one of its runs is
94/// still active.
95///
96/// # Examples
97///
98/// ```
99/// use ironflow_store::entities::OverlapPolicy;
100///
101/// assert_eq!(OverlapPolicy::default(), OverlapPolicy::Allow);
102/// assert_eq!(OverlapPolicy::Skip.as_ref(), "skip");
103///
104/// let parsed: OverlapPolicy = "allow".parse().unwrap();
105/// assert_eq!(parsed, OverlapPolicy::Allow);
106/// ```
107#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
108#[derive(
109 Debug,
110 Clone,
111 Copy,
112 Default,
113 PartialEq,
114 Eq,
115 Serialize,
116 Deserialize,
117 Display,
118 EnumString,
119 IntoStaticStr,
120 AsRefStr,
121)]
122#[serde(rename_all = "snake_case")]
123#[strum(serialize_all = "snake_case")]
124pub enum OverlapPolicy {
125 /// Create the run anyway: runs of the schedule may run side by side.
126 #[default]
127 Allow,
128 /// Skip the occurrence while a run of the schedule is active. Enforced
129 /// with the concurrency key [`Schedule::concurrency_key`].
130 Skip,
131}
132
133/// Default [`SchedulePolicy::catchup_max`].
134pub const DEFAULT_CATCHUP_MAX: u32 = 10;
135/// Lowest accepted [`SchedulePolicy::catchup_max`].
136pub const MIN_CATCHUP_MAX: u32 = 1;
137/// Highest accepted [`SchedulePolicy::catchup_max`].
138pub const MAX_CATCHUP_MAX: u32 = 1000;
139/// Default [`SchedulePolicy::catchup_window_secs`]: one day.
140pub const DEFAULT_CATCHUP_WINDOW_SECS: u32 = 86_400;
141/// Lowest accepted [`SchedulePolicy::catchup_window_secs`]: one minute.
142pub const MIN_CATCHUP_WINDOW_SECS: u32 = 60;
143/// Highest accepted [`SchedulePolicy::catchup_window_secs`]: thirty days.
144pub const MAX_CATCHUP_WINDOW_SECS: u32 = 2_592_000;
145/// Default [`SchedulePolicy::timezone`].
146pub const DEFAULT_TIMEZONE: Tz = Tz::UTC;
147
148/// Catch-up, overlap and timezone policy of a schedule.
149///
150/// # Examples
151///
152/// ```
153/// use chrono_tz::Tz;
154/// use ironflow_store::entities::{CatchupPolicy, OverlapPolicy, SchedulePolicy};
155///
156/// let policy = SchedulePolicy {
157/// catchup: CatchupPolicy::All,
158/// overlap: OverlapPolicy::Skip,
159/// timezone: Tz::Europe__Paris,
160/// ..SchedulePolicy::default()
161/// };
162/// assert!(policy.validate().is_ok());
163///
164/// let invalid = SchedulePolicy { catchup_max: 0, ..SchedulePolicy::default() };
165/// assert!(invalid.validate().is_err());
166/// ```
167#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
168#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
169pub struct SchedulePolicy {
170 /// What to do with missed occurrences.
171 pub catchup: CatchupPolicy,
172 /// Most runs created to catch up under [`CatchupPolicy::All`], between
173 /// [`MIN_CATCHUP_MAX`] and [`MAX_CATCHUP_MAX`].
174 pub catchup_max: u32,
175 /// How far back, in seconds, a missed occurrence is still caught up,
176 /// between [`MIN_CATCHUP_WINDOW_SECS`] and [`MAX_CATCHUP_WINDOW_SECS`].
177 pub catchup_window_secs: u32,
178 /// What to do when a run of the schedule is still active.
179 pub overlap: OverlapPolicy,
180 /// IANA timezone the cron expression is evaluated in, e.g. `Europe/Paris`.
181 #[cfg_attr(feature = "openapi", schema(value_type = String, example = "Europe/Paris"))]
182 pub timezone: Tz,
183}
184
185impl Default for SchedulePolicy {
186 fn default() -> Self {
187 Self {
188 catchup: CatchupPolicy::default(),
189 catchup_max: DEFAULT_CATCHUP_MAX,
190 catchup_window_secs: DEFAULT_CATCHUP_WINDOW_SECS,
191 overlap: OverlapPolicy::default(),
192 timezone: DEFAULT_TIMEZONE,
193 }
194 }
195}
196
197impl SchedulePolicy {
198 /// Check the numeric bounds of the policy. The timezone is checked where
199 /// the timezone database is available (API and engine).
200 ///
201 /// # Errors
202 ///
203 /// Returns a message naming the field out of its range.
204 ///
205 /// # Examples
206 ///
207 /// ```
208 /// use ironflow_store::entities::SchedulePolicy;
209 ///
210 /// let policy = SchedulePolicy { catchup_window_secs: 10, ..SchedulePolicy::default() };
211 /// assert_eq!(
212 /// policy.validate().unwrap_err(),
213 /// "catchup_window_secs must be between 60 and 2592000",
214 /// );
215 /// ```
216 pub fn validate(&self) -> Result<(), String> {
217 if !(MIN_CATCHUP_MAX..=MAX_CATCHUP_MAX).contains(&self.catchup_max) {
218 return Err(format!(
219 "catchup_max must be between {MIN_CATCHUP_MAX} and {MAX_CATCHUP_MAX}"
220 ));
221 }
222 if !(MIN_CATCHUP_WINDOW_SECS..=MAX_CATCHUP_WINDOW_SECS).contains(&self.catchup_window_secs)
223 {
224 return Err(format!(
225 "catchup_window_secs must be between {MIN_CATCHUP_WINDOW_SECS} and {MAX_CATCHUP_WINDOW_SECS}"
226 ));
227 }
228 Ok(())
229 }
230}
231
232/// Why an occurrence of a schedule created no run.
233///
234/// # Examples
235///
236/// ```
237/// use ironflow_store::entities::ScheduleMissReason;
238///
239/// let reason = ScheduleMissReason::Overlap;
240/// assert_eq!(reason.to_string(), "overlap");
241/// assert_eq!(serde_json::to_string(&reason).unwrap(), "\"overlap\"");
242/// ```
243#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
244#[derive(
245 Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, Display, EnumString, IntoStaticStr,
246)]
247#[serde(rename_all = "snake_case")]
248#[strum(serialize_all = "snake_case")]
249pub enum ScheduleMissReason {
250 /// Older than the schedule's catch-up window.
251 OutsideWindow,
252 /// Dropped because [`SchedulePolicy::catchup_max`] runs were already
253 /// caught up under [`CatchupPolicy::All`].
254 CatchupMax,
255 /// Replaced by a more recent occurrence under [`CatchupPolicy::Latest`].
256 Superseded,
257 /// Late, and not caught up under [`CatchupPolicy::Skip`].
258 CatchupSkip,
259 /// A run of the schedule was still active under [`OverlapPolicy::Skip`].
260 Overlap,
261}
262
263/// A persisted schedule that triggers a workflow on a cron expression.
264///
265/// A schedule is active when `disabled_at` is `None`. Setting `disabled_at`
266/// to a timestamp pauses it.
267///
268/// # Examples
269///
270/// ```
271/// use ironflow_store::entities::{Schedule, SchedulePolicy, ScheduleSource};
272/// use chrono::Utc;
273/// use serde_json::json;
274/// use uuid::Uuid;
275///
276/// let schedule = Schedule {
277/// id: Uuid::now_v7(),
278/// workflow_name: "deploy".to_string(),
279/// cron_expression: "0 0 * * * *".to_string(),
280/// inputs: json!({}),
281/// source: ScheduleSource::Api,
282/// disabled_at: None,
283/// last_triggered_at: None,
284/// next_trigger_at: Some(Utc::now()),
285/// last_error: None,
286/// priority: 0,
287/// policy: SchedulePolicy::default(),
288/// created_by_user_id: Some(Uuid::now_v7()),
289/// created_at: Utc::now(),
290/// updated_at: Utc::now(),
291/// };
292/// assert!(schedule.is_active());
293/// ```
294#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
295#[derive(Debug, Clone, Serialize, Deserialize)]
296pub struct Schedule {
297 /// Unique schedule ID (UUID v7).
298 pub id: Uuid,
299 /// Name of the workflow to trigger.
300 pub workflow_name: String,
301 /// Cron expression (6-field format, e.g. `"0 */5 * * * *"`).
302 pub cron_expression: String,
303 /// JSON payload passed to the workflow on each trigger.
304 pub inputs: Value,
305 /// Where this schedule was created.
306 pub source: ScheduleSource,
307 /// When the schedule was disabled. `None` means active.
308 pub disabled_at: Option<DateTime<Utc>>,
309 /// When the schedule last created a run.
310 pub last_triggered_at: Option<DateTime<Utc>>,
311 /// When the schedule will next fire. Always set on an active schedule.
312 pub next_trigger_at: Option<DateTime<Utc>>,
313 /// Why Ironflow disabled the schedule on its own, e.g. a cron expression
314 /// whose next occurrence cannot be computed. `None` for a schedule paused
315 /// by a user or never disabled.
316 #[serde(default)]
317 pub last_error: Option<String>,
318 /// User who created the schedule. `None` for handler-declared schedules,
319 /// which have no human author.
320 pub created_by_user_id: Option<Uuid>,
321 /// When the schedule was created.
322 pub created_at: DateTime<Utc>,
323 /// When the schedule was last updated.
324 pub updated_at: DateTime<Utc>,
325 /// Queue priority given to every run this schedule creates, between
326 /// [`MIN_PRIORITY`](super::MIN_PRIORITY) and [`MAX_PRIORITY`](super::MAX_PRIORITY).
327 ///
328 /// Defaults to `0` when absent from the payload.
329 #[serde(default)]
330 pub priority: i16,
331 /// Catch-up, overlap and timezone policy. Defaults to
332 /// [`SchedulePolicy::default`] when absent from the payload.
333 #[serde(default)]
334 pub policy: SchedulePolicy,
335}
336
337impl Schedule {
338 /// Whether the schedule is currently active (not disabled).
339 pub fn is_active(&self) -> bool {
340 self.disabled_at.is_none()
341 }
342
343 /// Idempotency key of the run created for one occurrence of this schedule:
344 /// `schedule:<id>:<occurrence in RFC 3339>`.
345 ///
346 /// Two firings of the same occurrence (several servers, a retry after a
347 /// crash) share the key, so they never create two runs.
348 ///
349 /// # Examples
350 ///
351 /// ```
352 /// use chrono::{TimeZone, Utc};
353 /// use ironflow_store::entities::Schedule;
354 /// use uuid::Uuid;
355 ///
356 /// let id = Uuid::nil();
357 /// let at = Utc.with_ymd_and_hms(2026, 10, 6, 8, 0, 0).unwrap();
358 /// assert_eq!(
359 /// Schedule::occurrence_key(id, at),
360 /// "schedule:00000000-0000-0000-0000-000000000000:2026-10-06T08:00:00Z",
361 /// );
362 /// ```
363 pub fn occurrence_key(id: Uuid, occurrence: DateTime<Utc>) -> String {
364 format!(
365 "schedule:{id}:{}",
366 occurrence.to_rfc3339_opts(SecondsFormat::AutoSi, true)
367 )
368 }
369
370 /// Concurrency key of the runs of a schedule under
371 /// [`OverlapPolicy::Skip`]: `schedule:<id>`.
372 ///
373 /// # Examples
374 ///
375 /// ```
376 /// use ironflow_store::entities::Schedule;
377 /// use uuid::Uuid;
378 ///
379 /// assert_eq!(
380 /// Schedule::concurrency_key(Uuid::nil()),
381 /// "schedule:00000000-0000-0000-0000-000000000000",
382 /// );
383 /// ```
384 pub fn concurrency_key(id: Uuid) -> String {
385 format!("schedule:{id}")
386 }
387
388 /// Build the run this schedule creates: its workflow, its inputs as
389 /// payload, and a [`TriggerKind::Cron`] trigger carrying the schedule id
390 /// and the occurrence it covers (`None` for a manual trigger).
391 ///
392 /// Under [`OverlapPolicy::Skip`] the run carries the concurrency key
393 /// [`Schedule::concurrency_key`], so it cannot start while another run
394 /// of the schedule is active.
395 ///
396 /// The run carries no idempotency key: the store sets one when it fires
397 /// an occurrence (see [`Schedule::occurrence_key`]).
398 ///
399 /// # Examples
400 ///
401 /// ```
402 /// use chrono::Utc;
403 /// use ironflow_store::entities::{Schedule, SchedulePolicy, ScheduleSource, TriggerKind};
404 /// use serde_json::json;
405 /// use uuid::Uuid;
406 ///
407 /// let schedule = Schedule {
408 /// id: Uuid::now_v7(),
409 /// workflow_name: "deploy".to_string(),
410 /// cron_expression: "0 0 * * * *".to_string(),
411 /// inputs: json!({"env": "prod"}),
412 /// source: ScheduleSource::Api,
413 /// disabled_at: None,
414 /// last_triggered_at: None,
415 /// next_trigger_at: Some(Utc::now()),
416 /// last_error: None,
417 /// priority: 0,
418 /// policy: SchedulePolicy::default(),
419 /// created_by_user_id: None,
420 /// created_at: Utc::now(),
421 /// updated_at: Utc::now(),
422 /// };
423 /// let occurrence = Utc::now();
424 /// let run = schedule.new_run(Some(occurrence), None);
425 /// assert_eq!(run.workflow_name, "deploy");
426 /// assert_eq!(run.payload, json!({"env": "prod"}));
427 /// assert!(matches!(
428 /// run.trigger,
429 /// TriggerKind::Cron { scheduled_for: Some(at), .. } if at == occurrence
430 /// ));
431 /// ```
432 pub fn new_run(
433 &self,
434 scheduled_for: Option<DateTime<Utc>>,
435 created_by: Option<RunActor>,
436 ) -> NewRun {
437 NewRun {
438 workflow_name: self.workflow_name.clone(),
439 trigger: TriggerKind::Cron {
440 schedule: self.cron_expression.clone(),
441 schedule_id: Some(self.id),
442 scheduled_for,
443 },
444 payload: self.inputs.clone(),
445 max_retries: 0,
446 handler_version: None,
447 labels: HashMap::new(),
448 scheduled_at: None,
449 created_by,
450 idempotency_key: None,
451 concurrency_key: (self.policy.overlap == OverlapPolicy::Skip)
452 .then(|| Schedule::concurrency_key(self.id)),
453 priority: self.priority,
454 concurrency_limits: Vec::new(),
455 max_cost_usd: None,
456 worker_tags: Vec::new(),
457 }
458 }
459}
460
461/// What happens to a schedule after it fires an occurrence.
462///
463/// Computed by the caller from the cron expression, applied by
464/// [`ScheduleStore::fire_due_schedule`](crate::schedule_store::ScheduleStore::fire_due_schedule)
465/// in the same transaction as the run creation.
466///
467/// # Examples
468///
469/// ```
470/// use chrono::Utc;
471/// use ironflow_store::entities::ScheduleNext;
472///
473/// let next = ScheduleNext::At(Utc::now());
474/// assert!(matches!(next, ScheduleNext::At(_)));
475///
476/// let stop = ScheduleNext::Disable { error: "no next occurrence".to_string() };
477/// assert!(matches!(stop, ScheduleNext::Disable { .. }));
478/// ```
479#[derive(Debug, Clone, PartialEq, Eq)]
480pub enum ScheduleNext {
481 /// Fire again at this time.
482 At(DateTime<Utc>),
483 /// Disable the schedule: its next occurrence cannot be computed. The
484 /// error is stored in [`Schedule::last_error`].
485 Disable {
486 /// Why the next occurrence cannot be computed.
487 error: String,
488 },
489}
490
491/// What one firing of a due schedule writes: the occurrences to create a run
492/// for, and what happens to the schedule afterwards.
493///
494/// Computed by the caller from the schedule's [`SchedulePolicy`], applied by
495/// [`ScheduleStore::fire_due_schedule`](crate::schedule_store::ScheduleStore::fire_due_schedule).
496///
497/// # Examples
498///
499/// ```
500/// use chrono::{TimeDelta, Utc};
501/// use ironflow_store::entities::{ScheduleFiringPlan, ScheduleNext};
502///
503/// let now = Utc::now();
504/// let plan = ScheduleFiringPlan {
505/// occurrences: vec![now - TimeDelta::hours(1), now],
506/// next: ScheduleNext::At(now + TimeDelta::hours(1)),
507/// };
508/// assert_eq!(plan.occurrences.len(), 2);
509/// ```
510#[derive(Debug, Clone, PartialEq, Eq)]
511pub struct ScheduleFiringPlan {
512 /// Occurrences to create a run for, oldest first. May be empty.
513 pub occurrences: Vec<DateTime<Utc>>,
514 /// What happens to the schedule after the firing.
515 pub next: ScheduleNext,
516}
517
518/// The run created, or found, for one occurrence of a schedule.
519///
520/// # Examples
521///
522/// ```no_run
523/// use ironflow_store::entities::ScheduledRun;
524///
525/// fn describe(run: &ScheduledRun) -> String {
526/// format!("{} covers {}", run.run.run().id, run.occurrence)
527/// }
528/// ```
529#[derive(Debug, Clone)]
530pub struct ScheduledRun {
531 /// The occurrence the run covers.
532 pub occurrence: DateTime<Utc>,
533 /// The run. [`RunCreation::Existing`] when a run with the same occurrence
534 /// key already existed.
535 pub run: RunCreation,
536}
537
538/// Result of firing a due schedule.
539///
540/// # Examples
541///
542/// ```no_run
543/// use chrono::Utc;
544/// use ironflow_store::entities::{ScheduleFiringPlan, ScheduleNext};
545/// use ironflow_store::memory::InMemoryStore;
546/// use ironflow_store::schedule_store::ScheduleStore;
547/// use uuid::Uuid;
548///
549/// # async fn example(id: Uuid, due: chrono::DateTime<Utc>) -> Result<(), ironflow_store::error::StoreError> {
550/// let store = InMemoryStore::new();
551/// let plan = ScheduleFiringPlan {
552/// occurrences: vec![due],
553/// next: ScheduleNext::At(Utc::now()),
554/// };
555/// if let Some(firing) = store.fire_due_schedule(id, due, plan).await? {
556/// for scheduled in &firing.runs {
557/// println!("run {} created for {}", scheduled.run.run().id, scheduled.occurrence);
558/// }
559/// println!("{} occurrences skipped as overlapping", firing.overlapped.len());
560/// }
561/// # Ok(())
562/// # }
563/// ```
564#[derive(Debug, Clone)]
565pub struct ScheduleFiring {
566 /// The schedule after the firing: next occurrence set, or disabled.
567 pub schedule: Schedule,
568 /// The runs of the fired occurrences, oldest first.
569 pub runs: Vec<ScheduledRun>,
570 /// Occurrences refused because a run holding the schedule's concurrency
571 /// key was still active ([`OverlapPolicy::Skip`]). Nothing was written
572 /// for them.
573 pub overlapped: Vec<DateTime<Utc>>,
574}
575
576/// Parameters for creating a new schedule.
577///
578/// # Examples
579///
580/// ```
581/// use ironflow_store::entities::{NewSchedule, SchedulePolicy, ScheduleSource};
582/// use serde_json::json;
583/// use uuid::Uuid;
584/// use chrono::Utc;
585///
586/// let new = NewSchedule {
587/// workflow_name: "deploy".to_string(),
588/// cron_expression: "0 0 * * * *".to_string(),
589/// inputs: json!({"env": "prod"}),
590/// source: ScheduleSource::Api,
591/// priority: 0,
592/// policy: SchedulePolicy::default(),
593/// created_by_user_id: Some(Uuid::now_v7()),
594/// next_trigger_at: Some(Utc::now()),
595/// };
596/// assert_eq!(new.workflow_name, "deploy");
597/// ```
598#[derive(Debug, Clone)]
599pub struct NewSchedule {
600 /// Name of the workflow to trigger.
601 pub workflow_name: String,
602 /// Cron expression (6-field format).
603 pub cron_expression: String,
604 /// JSON payload for the workflow.
605 pub inputs: Value,
606 /// Where this schedule originates.
607 pub source: ScheduleSource,
608 /// User who creates the schedule. `None` for handler-declared schedules,
609 /// which have no human author.
610 pub created_by_user_id: Option<Uuid>,
611 /// Pre-computed next trigger time.
612 pub next_trigger_at: Option<DateTime<Utc>>,
613 /// Queue priority given to every run the schedule creates, between
614 /// [`MIN_PRIORITY`](super::MIN_PRIORITY) and [`MAX_PRIORITY`](super::MAX_PRIORITY).
615 pub priority: i16,
616 /// Catch-up, overlap and timezone policy.
617 pub policy: SchedulePolicy,
618}
619
620/// Updatable fields on a schedule.
621///
622/// Only `Some` fields are applied; `None` means "leave unchanged".
623///
624/// # Examples
625///
626/// ```
627/// use ironflow_store::entities::ScheduleUpdate;
628/// use serde_json::json;
629///
630/// let update = ScheduleUpdate {
631/// cron_expression: Some("0 30 * * * *".to_string()),
632/// inputs: Some(json!({"env": "staging"})),
633/// disabled_at: None,
634/// next_trigger_at: None,
635/// last_triggered_at: None,
636/// priority: None,
637/// last_error: None,
638/// policy: None,
639/// };
640/// assert!(update.disabled_at.is_none());
641/// ```
642#[derive(Debug, Clone, Default)]
643pub struct ScheduleUpdate {
644 /// New cron expression.
645 pub cron_expression: Option<String>,
646 /// New inputs payload.
647 pub inputs: Option<Value>,
648 /// Set or clear disabled_at. `Some(Some(ts))` disables, `Some(None)` re-enables.
649 pub disabled_at: Option<Option<DateTime<Utc>>>,
650 /// Updated next trigger time.
651 pub next_trigger_at: Option<Option<DateTime<Utc>>>,
652 /// Updated last triggered time.
653 pub last_triggered_at: Option<Option<DateTime<Utc>>>,
654 /// Set or clear [`Schedule::last_error`]. `Some(None)` clears it.
655 pub last_error: Option<Option<String>>,
656 /// New queue priority for the runs the schedule creates.
657 pub priority: Option<i16>,
658 /// New catch-up, overlap and timezone policy.
659 pub policy: Option<SchedulePolicy>,
660}
661
662#[cfg(test)]
663mod tests {
664 use super::*;
665 use serde_json::json;
666
667 fn schedule(policy: SchedulePolicy) -> Schedule {
668 Schedule {
669 id: Uuid::now_v7(),
670 workflow_name: "deploy".to_string(),
671 cron_expression: "0 0 * * * *".to_string(),
672 inputs: json!({}),
673 source: ScheduleSource::Api,
674 disabled_at: None,
675 last_triggered_at: None,
676 next_trigger_at: Some(Utc::now()),
677 last_error: None,
678 priority: 42,
679 policy,
680 created_by_user_id: None,
681 created_at: Utc::now(),
682 updated_at: Utc::now(),
683 }
684 }
685
686 #[test]
687 fn schedule_serde_roundtrip() {
688 let schedule = Schedule {
689 id: Uuid::now_v7(),
690 workflow_name: "deploy".to_string(),
691 cron_expression: "0 0 * * * *".to_string(),
692 inputs: json!({"env": "prod"}),
693 source: ScheduleSource::Api,
694 disabled_at: None,
695 last_triggered_at: None,
696 next_trigger_at: Some(Utc::now()),
697 last_error: None,
698 priority: -5,
699 policy: SchedulePolicy {
700 catchup: CatchupPolicy::All,
701 timezone: Tz::Europe__Paris,
702 ..SchedulePolicy::default()
703 },
704 created_by_user_id: Some(Uuid::now_v7()),
705 created_at: Utc::now(),
706 updated_at: Utc::now(),
707 };
708 let json_str = serde_json::to_string(&schedule).expect("serialize");
709 let back: Schedule = serde_json::from_str(&json_str).expect("deserialize");
710 assert_eq!(schedule.id, back.id);
711 assert_eq!(schedule.workflow_name, back.workflow_name);
712 assert_eq!(back.priority, -5);
713 assert_eq!(back.policy, schedule.policy);
714 assert!(back.is_active());
715 }
716
717 #[test]
718 fn schedule_without_policy_deserializes_with_defaults() {
719 let mut value = serde_json::to_value(schedule(SchedulePolicy::default())).unwrap();
720 value.as_object_mut().unwrap().remove("policy");
721 let back: Schedule = serde_json::from_value(value).expect("deserialize");
722 assert_eq!(back.policy, SchedulePolicy::default());
723 }
724
725 #[test]
726 fn schedule_new_run_carries_the_schedule_priority() {
727 let schedule = schedule(SchedulePolicy::default());
728 assert_eq!(schedule.new_run(None, None).priority, 42);
729 }
730
731 #[test]
732 fn new_run_carries_schedule_id_and_occurrence() {
733 let schedule = schedule(SchedulePolicy::default());
734 let occurrence = Utc::now();
735 let run = schedule.new_run(Some(occurrence), None);
736 assert_eq!(
737 run.trigger,
738 TriggerKind::Cron {
739 schedule: "0 0 * * * *".to_string(),
740 schedule_id: Some(schedule.id),
741 scheduled_for: Some(occurrence),
742 }
743 );
744 assert!(run.concurrency_key.is_none());
745
746 let manual = schedule.new_run(None, None);
747 assert!(matches!(
748 manual.trigger,
749 TriggerKind::Cron { schedule_id: Some(id), scheduled_for: None, .. } if id == schedule.id
750 ));
751 }
752
753 #[test]
754 fn new_run_with_overlap_skip_sets_the_schedule_concurrency_key() {
755 let schedule = schedule(SchedulePolicy {
756 overlap: OverlapPolicy::Skip,
757 ..SchedulePolicy::default()
758 });
759 let run = schedule.new_run(Some(Utc::now()), None);
760 assert_eq!(
761 run.concurrency_key,
762 Some(Schedule::concurrency_key(schedule.id))
763 );
764 }
765
766 #[test]
767 fn disabled_schedule_is_not_active() {
768 let schedule = Schedule {
769 id: Uuid::now_v7(),
770 workflow_name: "deploy".to_string(),
771 cron_expression: "0 0 * * * *".to_string(),
772 inputs: json!({}),
773 source: ScheduleSource::Api,
774 disabled_at: Some(Utc::now()),
775 last_triggered_at: None,
776 next_trigger_at: None,
777 last_error: None,
778 priority: 0,
779 policy: SchedulePolicy::default(),
780 created_by_user_id: Some(Uuid::now_v7()),
781 created_at: Utc::now(),
782 updated_at: Utc::now(),
783 };
784 assert!(!schedule.is_active());
785 }
786
787 #[test]
788 fn schedule_source_roundtrip() {
789 assert_eq!(ScheduleSource::Handler.as_str(), "handler");
790 assert_eq!(ScheduleSource::Api.as_str(), "api");
791
792 let parsed: ScheduleSource = "handler".parse().unwrap();
793 assert_eq!(parsed, ScheduleSource::Handler);
794
795 let parsed: ScheduleSource = "api".parse().unwrap();
796 assert_eq!(parsed, ScheduleSource::Api);
797
798 assert!("unknown".parse::<ScheduleSource>().is_err());
799 }
800
801 #[test]
802 fn catchup_and_overlap_policies_roundtrip() {
803 for policy in [
804 CatchupPolicy::Latest,
805 CatchupPolicy::All,
806 CatchupPolicy::Skip,
807 ] {
808 assert_eq!(policy.as_ref().parse::<CatchupPolicy>().unwrap(), policy);
809 }
810 for policy in [OverlapPolicy::Allow, OverlapPolicy::Skip] {
811 assert_eq!(policy.as_ref().parse::<OverlapPolicy>().unwrap(), policy);
812 }
813 assert!("never".parse::<CatchupPolicy>().is_err());
814 assert!("queue".parse::<OverlapPolicy>().is_err());
815 }
816
817 #[test]
818 fn schedule_policy_defaults() {
819 let policy = SchedulePolicy::default();
820 assert_eq!(policy.catchup, CatchupPolicy::Latest);
821 assert_eq!(policy.catchup_max, DEFAULT_CATCHUP_MAX);
822 assert_eq!(policy.catchup_window_secs, DEFAULT_CATCHUP_WINDOW_SECS);
823 assert_eq!(policy.overlap, OverlapPolicy::Allow);
824 assert_eq!(policy.timezone, Tz::UTC);
825 assert!(policy.validate().is_ok());
826 }
827
828 #[test]
829 fn schedule_policy_validate_checks_bounds() {
830 let at = |catchup_max, catchup_window_secs| SchedulePolicy {
831 catchup_max,
832 catchup_window_secs,
833 ..SchedulePolicy::default()
834 };
835 assert!(
836 at(MIN_CATCHUP_MAX, MIN_CATCHUP_WINDOW_SECS)
837 .validate()
838 .is_ok()
839 );
840 assert!(
841 at(MAX_CATCHUP_MAX, MAX_CATCHUP_WINDOW_SECS)
842 .validate()
843 .is_ok()
844 );
845 assert_eq!(
846 at(0, 3600).validate().unwrap_err(),
847 "catchup_max must be between 1 and 1000"
848 );
849 assert_eq!(
850 at(1001, 3600).validate().unwrap_err(),
851 "catchup_max must be between 1 and 1000"
852 );
853 assert_eq!(
854 at(10, 59).validate().unwrap_err(),
855 "catchup_window_secs must be between 60 and 2592000"
856 );
857 assert!(at(10, MAX_CATCHUP_WINDOW_SECS + 1).validate().is_err());
858 }
859
860 #[test]
861 fn schedule_miss_reason_serializes_snake_case() {
862 assert_eq!(
863 serde_json::to_string(&ScheduleMissReason::OutsideWindow).unwrap(),
864 "\"outside_window\""
865 );
866 assert_eq!(ScheduleMissReason::CatchupMax.to_string(), "catchup_max");
867 assert_eq!(ScheduleMissReason::CatchupSkip.to_string(), "catchup_skip");
868 }
869
870 #[test]
871 fn schedule_update_defaults_to_none() {
872 let update = ScheduleUpdate::default();
873 assert!(update.cron_expression.is_none());
874 assert!(update.inputs.is_none());
875 assert!(update.disabled_at.is_none());
876 assert!(update.next_trigger_at.is_none());
877 assert!(update.last_triggered_at.is_none());
878 assert!(update.last_error.is_none());
879 assert!(update.priority.is_none());
880 assert!(update.policy.is_none());
881 }
882}