1use std::collections::HashMap;
11use std::sync::Arc;
12use std::time::Duration;
13
14use croner::Cron;
15use ironflow_store::entities::{NewRun, RunActor, Schedule, ScheduleUpdate, TriggerKind};
16use ironflow_store::store::Store;
17use tokio::time::interval;
18use tokio_util::sync::CancellationToken;
19use tracing::{error, info, warn};
20
21pub(crate) fn next_trigger(
23 cron_str: &str,
24) -> Result<Option<chrono::DateTime<chrono::Utc>>, String> {
25 let mut cron = Cron::new(cron_str);
26 cron.pattern.with_seconds_optional = true;
27 let cron = cron
28 .parse()
29 .map_err(|e| format!("invalid cron expression: {e}"))?;
30 let next = cron
31 .find_next_occurrence(&chrono::Utc::now(), false)
32 .map_err(|e| format!("cannot compute next trigger: {e}"))?;
33 Ok(Some(next))
34}
35
36pub(crate) fn new_run_from_schedule(schedule: &Schedule, created_by: Option<RunActor>) -> NewRun {
38 NewRun {
39 workflow_name: schedule.workflow_name.clone(),
40 trigger: TriggerKind::Cron {
41 schedule: schedule.cron_expression.clone(),
42 },
43 payload: schedule.inputs.clone(),
44 max_retries: 0,
45 handler_version: None,
46 labels: HashMap::new(),
47 scheduled_at: None,
48 created_by,
49 idempotency_key: None,
50 max_cost_usd: None,
51 }
52}
53
54pub const DEFAULT_TICK_INTERVAL: Duration = Duration::from_secs(15);
56
57pub struct ScheduleTicker {
83 store: Arc<dyn Store>,
84 interval: Duration,
85}
86
87impl ScheduleTicker {
88 pub fn new(store: Arc<dyn Store>) -> Self {
90 Self {
91 store,
92 interval: DEFAULT_TICK_INTERVAL,
93 }
94 }
95
96 pub fn interval(mut self, interval: Duration) -> Self {
98 self.interval = interval;
99 self
100 }
101
102 pub async fn run(self, shutdown: CancellationToken) {
106 let mut ticker = interval(self.interval);
107 ticker.tick().await;
108
109 info!(
110 interval_secs = self.interval.as_secs(),
111 "schedule ticker started"
112 );
113
114 loop {
115 tokio::select! {
116 _ = shutdown.cancelled() => {
117 info!("schedule ticker stopped");
118 return;
119 }
120 _ = ticker.tick() => {
121 self.tick().await;
122 }
123 }
124 }
125 }
126
127 pub async fn tick(&self) {
131 let claimed = match self.store.claim_due_schedules().await {
132 Ok(claimed) => claimed,
133 Err(err) => {
134 error!(error = %err, "failed to claim due schedules");
135 return;
136 }
137 };
138
139 if claimed.is_empty() {
140 return;
141 }
142
143 info!(count = claimed.len(), "firing due schedules");
144
145 for schedule in claimed {
146 let run_result = self
147 .store
148 .create_run(new_run_from_schedule(&schedule, None))
149 .await;
150
151 match run_result {
152 Ok(creation) => {
153 let run = creation.into_run();
154 info!(
155 schedule_id = %schedule.id,
156 workflow = %schedule.workflow_name,
157 run_id = %run.id,
158 "schedule fired"
159 );
160 }
161 Err(err) => {
162 error!(
163 schedule_id = %schedule.id,
164 workflow = %schedule.workflow_name,
165 error = %err,
166 "failed to create run for schedule"
167 );
168 continue;
169 }
170 }
171
172 let next = match next_trigger(&schedule.cron_expression) {
173 Ok(next) => next,
174 Err(err) => {
175 warn!(
176 schedule_id = %schedule.id,
177 error = %err,
178 "cannot compute next trigger, schedule will remain paused"
179 );
180 None
181 }
182 };
183
184 if let Err(err) = self
185 .store
186 .update_schedule(
187 schedule.id,
188 ScheduleUpdate {
189 next_trigger_at: Some(next),
190 ..Default::default()
191 },
192 )
193 .await
194 {
195 error!(
196 schedule_id = %schedule.id,
197 error = %err,
198 "failed to set next trigger time after firing"
199 );
200 }
201 }
202 }
203}
204
205#[cfg(test)]
206mod tests {
207 use chrono::Utc;
208 use ironflow_store::entities::{NewSchedule, RunFilter, ScheduleSource};
209 use ironflow_store::memory::InMemoryStore;
210 use ironflow_store::store::Store;
211 use serde_json::json;
212 use std::sync::Arc;
213 use uuid::Uuid;
214
215 use super::*;
216
217 async fn make_store_with_due_schedule() -> (Arc<dyn Store>, Uuid) {
218 let store: Arc<dyn Store> = Arc::new(InMemoryStore::new());
219 let schedule = store
220 .create_schedule(NewSchedule {
221 workflow_name: "deploy".to_string(),
222 cron_expression: "* * * * *".to_string(),
223 inputs: json!({}),
224 source: ScheduleSource::Api,
225 created_by_user_id: Uuid::now_v7(),
226 next_trigger_at: Some(Utc::now() - chrono::Duration::seconds(10)),
227 })
228 .await
229 .expect("create schedule");
230 (store, schedule.id)
231 }
232
233 #[tokio::test]
234 async fn tick_creates_run_for_due_schedule() {
235 let (store, schedule_id) = make_store_with_due_schedule().await;
236 let ticker = ScheduleTicker::new(store.clone());
237
238 ticker.tick().await;
239
240 let runs = store
241 .list_runs(RunFilter::default(), 1, 10)
242 .await
243 .expect("list runs");
244 assert_eq!(runs.items.len(), 1);
245 assert_eq!(runs.items[0].workflow_name, "deploy");
246
247 let updated = store
248 .find_schedule_by_id(schedule_id)
249 .await
250 .expect("find")
251 .expect("exists");
252 assert!(updated.last_triggered_at.is_some());
253 assert!(updated.next_trigger_at.is_some());
254 assert!(updated.next_trigger_at.unwrap() > Utc::now() - chrono::Duration::seconds(1));
255 }
256
257 #[tokio::test]
258 async fn tick_skips_disabled_schedule() {
259 let store: Arc<dyn Store> = Arc::new(InMemoryStore::new());
260 let schedule = store
261 .create_schedule(NewSchedule {
262 workflow_name: "deploy".to_string(),
263 cron_expression: "* * * * *".to_string(),
264 inputs: json!({}),
265 source: ScheduleSource::Api,
266 created_by_user_id: Uuid::now_v7(),
267 next_trigger_at: Some(Utc::now() - chrono::Duration::seconds(10)),
268 })
269 .await
270 .expect("create");
271
272 store
273 .update_schedule(
274 schedule.id,
275 ScheduleUpdate {
276 disabled_at: Some(Some(Utc::now())),
277 ..Default::default()
278 },
279 )
280 .await
281 .expect("disable");
282
283 let ticker = ScheduleTicker::new(store.clone());
284 ticker.tick().await;
285
286 let runs = store
287 .list_runs(RunFilter::default(), 1, 10)
288 .await
289 .expect("list runs");
290 assert!(runs.items.is_empty());
291 }
292
293 #[tokio::test]
294 async fn tick_skips_future_schedule() {
295 let store: Arc<dyn Store> = Arc::new(InMemoryStore::new());
296 store
297 .create_schedule(NewSchedule {
298 workflow_name: "deploy".to_string(),
299 cron_expression: "* * * * *".to_string(),
300 inputs: json!({}),
301 source: ScheduleSource::Api,
302 created_by_user_id: Uuid::now_v7(),
303 next_trigger_at: Some(Utc::now() + chrono::Duration::hours(1)),
304 })
305 .await
306 .expect("create");
307
308 let ticker = ScheduleTicker::new(store.clone());
309 ticker.tick().await;
310
311 let runs = store
312 .list_runs(RunFilter::default(), 1, 10)
313 .await
314 .expect("list runs");
315 assert!(runs.items.is_empty());
316 }
317
318 #[tokio::test]
319 async fn tick_does_nothing_when_empty() {
320 let store: Arc<dyn Store> = Arc::new(InMemoryStore::new());
321 let ticker = ScheduleTicker::new(store.clone());
322 ticker.tick().await;
323
324 let runs = store
325 .list_runs(RunFilter::default(), 1, 10)
326 .await
327 .expect("list runs");
328 assert!(runs.items.is_empty());
329 }
330
331 #[tokio::test]
332 async fn tick_does_not_double_fire() {
333 let (store, _) = make_store_with_due_schedule().await;
334 let ticker = ScheduleTicker::new(store.clone());
335
336 ticker.tick().await;
337 ticker.tick().await;
338
339 let runs = store
340 .list_runs(RunFilter::default(), 1, 10)
341 .await
342 .expect("list runs");
343 assert_eq!(
344 runs.items.len(),
345 1,
346 "second tick must not create a duplicate run"
347 );
348 }
349}