use std::collections::HashMap;
use chrono::{DateTime, SecondsFormat, Utc};
use chrono_tz::Tz;
use serde::{Deserialize, Serialize};
use serde_json::Value;
use strum::{AsRefStr, Display, EnumString, IntoStaticStr};
use uuid::Uuid;
use super::{NewRun, RunActor, RunCreation, TriggerKind};
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
#[derive(
Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, Display, EnumString, IntoStaticStr,
)]
#[serde(rename_all = "snake_case")]
#[strum(serialize_all = "snake_case")]
pub enum ScheduleSource {
Handler,
Api,
}
impl ScheduleSource {
pub fn as_str(&self) -> &'static str {
self.into()
}
}
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
#[derive(
Debug,
Clone,
Copy,
Default,
PartialEq,
Eq,
Serialize,
Deserialize,
Display,
EnumString,
IntoStaticStr,
AsRefStr,
)]
#[serde(rename_all = "snake_case")]
#[strum(serialize_all = "snake_case")]
pub enum CatchupPolicy {
#[default]
Latest,
All,
Skip,
}
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
#[derive(
Debug,
Clone,
Copy,
Default,
PartialEq,
Eq,
Serialize,
Deserialize,
Display,
EnumString,
IntoStaticStr,
AsRefStr,
)]
#[serde(rename_all = "snake_case")]
#[strum(serialize_all = "snake_case")]
pub enum OverlapPolicy {
#[default]
Allow,
Skip,
}
pub const DEFAULT_CATCHUP_MAX: u32 = 10;
pub const MIN_CATCHUP_MAX: u32 = 1;
pub const MAX_CATCHUP_MAX: u32 = 1000;
pub const DEFAULT_CATCHUP_WINDOW_SECS: u32 = 86_400;
pub const MIN_CATCHUP_WINDOW_SECS: u32 = 60;
pub const MAX_CATCHUP_WINDOW_SECS: u32 = 2_592_000;
pub const DEFAULT_TIMEZONE: Tz = Tz::UTC;
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct SchedulePolicy {
pub catchup: CatchupPolicy,
pub catchup_max: u32,
pub catchup_window_secs: u32,
pub overlap: OverlapPolicy,
#[cfg_attr(feature = "openapi", schema(value_type = String, example = "Europe/Paris"))]
pub timezone: Tz,
}
impl Default for SchedulePolicy {
fn default() -> Self {
Self {
catchup: CatchupPolicy::default(),
catchup_max: DEFAULT_CATCHUP_MAX,
catchup_window_secs: DEFAULT_CATCHUP_WINDOW_SECS,
overlap: OverlapPolicy::default(),
timezone: DEFAULT_TIMEZONE,
}
}
}
impl SchedulePolicy {
pub fn validate(&self) -> Result<(), String> {
if !(MIN_CATCHUP_MAX..=MAX_CATCHUP_MAX).contains(&self.catchup_max) {
return Err(format!(
"catchup_max must be between {MIN_CATCHUP_MAX} and {MAX_CATCHUP_MAX}"
));
}
if !(MIN_CATCHUP_WINDOW_SECS..=MAX_CATCHUP_WINDOW_SECS).contains(&self.catchup_window_secs)
{
return Err(format!(
"catchup_window_secs must be between {MIN_CATCHUP_WINDOW_SECS} and {MAX_CATCHUP_WINDOW_SECS}"
));
}
Ok(())
}
}
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
#[derive(
Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, Display, EnumString, IntoStaticStr,
)]
#[serde(rename_all = "snake_case")]
#[strum(serialize_all = "snake_case")]
pub enum ScheduleMissReason {
OutsideWindow,
CatchupMax,
Superseded,
CatchupSkip,
Overlap,
}
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct Schedule {
pub id: Uuid,
pub workflow_name: String,
pub cron_expression: String,
pub inputs: Value,
pub source: ScheduleSource,
pub disabled_at: Option<DateTime<Utc>>,
pub last_triggered_at: Option<DateTime<Utc>>,
pub next_trigger_at: Option<DateTime<Utc>>,
#[serde(default)]
pub last_error: Option<String>,
pub created_by_user_id: Option<Uuid>,
pub created_at: DateTime<Utc>,
pub updated_at: DateTime<Utc>,
#[serde(default)]
pub priority: i16,
#[serde(default)]
pub policy: SchedulePolicy,
}
impl Schedule {
pub fn is_active(&self) -> bool {
self.disabled_at.is_none()
}
pub fn occurrence_key(id: Uuid, occurrence: DateTime<Utc>) -> String {
format!(
"schedule:{id}:{}",
occurrence.to_rfc3339_opts(SecondsFormat::AutoSi, true)
)
}
pub fn concurrency_key(id: Uuid) -> String {
format!("schedule:{id}")
}
pub fn new_run(
&self,
scheduled_for: Option<DateTime<Utc>>,
created_by: Option<RunActor>,
) -> NewRun {
NewRun {
workflow_name: self.workflow_name.clone(),
trigger: TriggerKind::Cron {
schedule: self.cron_expression.clone(),
schedule_id: Some(self.id),
scheduled_for,
},
payload: self.inputs.clone(),
max_retries: 0,
handler_version: None,
labels: HashMap::new(),
scheduled_at: None,
created_by,
idempotency_key: None,
concurrency_key: (self.policy.overlap == OverlapPolicy::Skip)
.then(|| Schedule::concurrency_key(self.id)),
priority: self.priority,
concurrency_limits: Vec::new(),
max_cost_usd: None,
worker_tags: Vec::new(),
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ScheduleNext {
At(DateTime<Utc>),
Disable {
error: String,
},
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ScheduleFiringPlan {
pub occurrences: Vec<DateTime<Utc>>,
pub next: ScheduleNext,
}
#[derive(Debug, Clone)]
pub struct ScheduledRun {
pub occurrence: DateTime<Utc>,
pub run: RunCreation,
}
#[derive(Debug, Clone)]
pub struct ScheduleFiring {
pub schedule: Schedule,
pub runs: Vec<ScheduledRun>,
pub overlapped: Vec<DateTime<Utc>>,
}
#[derive(Debug, Clone)]
pub struct NewSchedule {
pub workflow_name: String,
pub cron_expression: String,
pub inputs: Value,
pub source: ScheduleSource,
pub created_by_user_id: Option<Uuid>,
pub next_trigger_at: Option<DateTime<Utc>>,
pub priority: i16,
pub policy: SchedulePolicy,
}
#[derive(Debug, Clone, Default)]
pub struct ScheduleUpdate {
pub cron_expression: Option<String>,
pub inputs: Option<Value>,
pub disabled_at: Option<Option<DateTime<Utc>>>,
pub next_trigger_at: Option<Option<DateTime<Utc>>>,
pub last_triggered_at: Option<Option<DateTime<Utc>>>,
pub last_error: Option<Option<String>>,
pub priority: Option<i16>,
pub policy: Option<SchedulePolicy>,
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
fn schedule(policy: SchedulePolicy) -> Schedule {
Schedule {
id: Uuid::now_v7(),
workflow_name: "deploy".to_string(),
cron_expression: "0 0 * * * *".to_string(),
inputs: json!({}),
source: ScheduleSource::Api,
disabled_at: None,
last_triggered_at: None,
next_trigger_at: Some(Utc::now()),
last_error: None,
priority: 42,
policy,
created_by_user_id: None,
created_at: Utc::now(),
updated_at: Utc::now(),
}
}
#[test]
fn schedule_serde_roundtrip() {
let schedule = Schedule {
id: Uuid::now_v7(),
workflow_name: "deploy".to_string(),
cron_expression: "0 0 * * * *".to_string(),
inputs: json!({"env": "prod"}),
source: ScheduleSource::Api,
disabled_at: None,
last_triggered_at: None,
next_trigger_at: Some(Utc::now()),
last_error: None,
priority: -5,
policy: SchedulePolicy {
catchup: CatchupPolicy::All,
timezone: Tz::Europe__Paris,
..SchedulePolicy::default()
},
created_by_user_id: Some(Uuid::now_v7()),
created_at: Utc::now(),
updated_at: Utc::now(),
};
let json_str = serde_json::to_string(&schedule).expect("serialize");
let back: Schedule = serde_json::from_str(&json_str).expect("deserialize");
assert_eq!(schedule.id, back.id);
assert_eq!(schedule.workflow_name, back.workflow_name);
assert_eq!(back.priority, -5);
assert_eq!(back.policy, schedule.policy);
assert!(back.is_active());
}
#[test]
fn schedule_without_policy_deserializes_with_defaults() {
let mut value = serde_json::to_value(schedule(SchedulePolicy::default())).unwrap();
value.as_object_mut().unwrap().remove("policy");
let back: Schedule = serde_json::from_value(value).expect("deserialize");
assert_eq!(back.policy, SchedulePolicy::default());
}
#[test]
fn schedule_new_run_carries_the_schedule_priority() {
let schedule = schedule(SchedulePolicy::default());
assert_eq!(schedule.new_run(None, None).priority, 42);
}
#[test]
fn new_run_carries_schedule_id_and_occurrence() {
let schedule = schedule(SchedulePolicy::default());
let occurrence = Utc::now();
let run = schedule.new_run(Some(occurrence), None);
assert_eq!(
run.trigger,
TriggerKind::Cron {
schedule: "0 0 * * * *".to_string(),
schedule_id: Some(schedule.id),
scheduled_for: Some(occurrence),
}
);
assert!(run.concurrency_key.is_none());
let manual = schedule.new_run(None, None);
assert!(matches!(
manual.trigger,
TriggerKind::Cron { schedule_id: Some(id), scheduled_for: None, .. } if id == schedule.id
));
}
#[test]
fn new_run_with_overlap_skip_sets_the_schedule_concurrency_key() {
let schedule = schedule(SchedulePolicy {
overlap: OverlapPolicy::Skip,
..SchedulePolicy::default()
});
let run = schedule.new_run(Some(Utc::now()), None);
assert_eq!(
run.concurrency_key,
Some(Schedule::concurrency_key(schedule.id))
);
}
#[test]
fn disabled_schedule_is_not_active() {
let schedule = Schedule {
id: Uuid::now_v7(),
workflow_name: "deploy".to_string(),
cron_expression: "0 0 * * * *".to_string(),
inputs: json!({}),
source: ScheduleSource::Api,
disabled_at: Some(Utc::now()),
last_triggered_at: None,
next_trigger_at: None,
last_error: None,
priority: 0,
policy: SchedulePolicy::default(),
created_by_user_id: Some(Uuid::now_v7()),
created_at: Utc::now(),
updated_at: Utc::now(),
};
assert!(!schedule.is_active());
}
#[test]
fn schedule_source_roundtrip() {
assert_eq!(ScheduleSource::Handler.as_str(), "handler");
assert_eq!(ScheduleSource::Api.as_str(), "api");
let parsed: ScheduleSource = "handler".parse().unwrap();
assert_eq!(parsed, ScheduleSource::Handler);
let parsed: ScheduleSource = "api".parse().unwrap();
assert_eq!(parsed, ScheduleSource::Api);
assert!("unknown".parse::<ScheduleSource>().is_err());
}
#[test]
fn catchup_and_overlap_policies_roundtrip() {
for policy in [
CatchupPolicy::Latest,
CatchupPolicy::All,
CatchupPolicy::Skip,
] {
assert_eq!(policy.as_ref().parse::<CatchupPolicy>().unwrap(), policy);
}
for policy in [OverlapPolicy::Allow, OverlapPolicy::Skip] {
assert_eq!(policy.as_ref().parse::<OverlapPolicy>().unwrap(), policy);
}
assert!("never".parse::<CatchupPolicy>().is_err());
assert!("queue".parse::<OverlapPolicy>().is_err());
}
#[test]
fn schedule_policy_defaults() {
let policy = SchedulePolicy::default();
assert_eq!(policy.catchup, CatchupPolicy::Latest);
assert_eq!(policy.catchup_max, DEFAULT_CATCHUP_MAX);
assert_eq!(policy.catchup_window_secs, DEFAULT_CATCHUP_WINDOW_SECS);
assert_eq!(policy.overlap, OverlapPolicy::Allow);
assert_eq!(policy.timezone, Tz::UTC);
assert!(policy.validate().is_ok());
}
#[test]
fn schedule_policy_validate_checks_bounds() {
let at = |catchup_max, catchup_window_secs| SchedulePolicy {
catchup_max,
catchup_window_secs,
..SchedulePolicy::default()
};
assert!(
at(MIN_CATCHUP_MAX, MIN_CATCHUP_WINDOW_SECS)
.validate()
.is_ok()
);
assert!(
at(MAX_CATCHUP_MAX, MAX_CATCHUP_WINDOW_SECS)
.validate()
.is_ok()
);
assert_eq!(
at(0, 3600).validate().unwrap_err(),
"catchup_max must be between 1 and 1000"
);
assert_eq!(
at(1001, 3600).validate().unwrap_err(),
"catchup_max must be between 1 and 1000"
);
assert_eq!(
at(10, 59).validate().unwrap_err(),
"catchup_window_secs must be between 60 and 2592000"
);
assert!(at(10, MAX_CATCHUP_WINDOW_SECS + 1).validate().is_err());
}
#[test]
fn schedule_miss_reason_serializes_snake_case() {
assert_eq!(
serde_json::to_string(&ScheduleMissReason::OutsideWindow).unwrap(),
"\"outside_window\""
);
assert_eq!(ScheduleMissReason::CatchupMax.to_string(), "catchup_max");
assert_eq!(ScheduleMissReason::CatchupSkip.to_string(), "catchup_skip");
}
#[test]
fn schedule_update_defaults_to_none() {
let update = ScheduleUpdate::default();
assert!(update.cron_expression.is_none());
assert!(update.inputs.is_none());
assert!(update.disabled_at.is_none());
assert!(update.next_trigger_at.is_none());
assert!(update.last_triggered_at.is_none());
assert!(update.last_error.is_none());
assert!(update.priority.is_none());
assert!(update.policy.is_none());
}
}