Skip to main content

ironflow_api/
schedule_ticker.rs

1//! Unified schedule executor for all workflow schedules.
2//!
3//! All schedules live in the database -- both those created via the REST API
4//! and those declared by a [`WorkflowHandler::schedule()`]. At server startup,
5//! call [`sync_handler_schedules`](crate::schedule_sync::sync_handler_schedules)
6//! to reconcile handler-declared schedules, then spawn [`ScheduleTicker::run`]
7//! which polls [`claim_due_schedules`] and creates a run for each claimed
8//! schedule.
9
10use 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
21/// Compute the next trigger time from a cron expression (5-field standard format).
22pub(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
36/// Build a [`NewRun`] from a schedule's fields.
37pub(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
54/// How often the ticker checks for due schedules.
55pub const DEFAULT_TICK_INTERVAL: Duration = Duration::from_secs(15);
56
57/// Periodic task that fires DB-persisted schedules.
58///
59/// # Examples
60///
61/// ```no_run
62/// use std::sync::Arc;
63/// use std::time::Duration;
64/// use ironflow_api::schedule_ticker::ScheduleTicker;
65/// use ironflow_api::schedule_sync::sync_handler_schedules;
66/// use ironflow_engine::engine::Engine;
67/// use ironflow_core::providers::claude::ClaudeCodeProvider;
68/// use ironflow_store::memory::InMemoryStore;
69/// use ironflow_store::store::Store;
70/// use tokio_util::sync::CancellationToken;
71///
72/// # async fn example() -> Result<(), Box<dyn std::error::Error>> {
73/// let store: Arc<dyn Store> = Arc::new(InMemoryStore::new());
74/// let engine = Engine::new(store.clone(), Arc::new(ClaudeCodeProvider::new()));
75/// sync_handler_schedules(&engine, store.as_ref()).await?;
76///
77/// let ticker = ScheduleTicker::new(store).interval(Duration::from_secs(10));
78/// tokio::spawn(ticker.run(CancellationToken::new()));
79/// # Ok(())
80/// # }
81/// ```
82pub struct ScheduleTicker {
83    store: Arc<dyn Store>,
84    interval: Duration,
85}
86
87impl ScheduleTicker {
88    /// Create a ticker with the default interval.
89    pub fn new(store: Arc<dyn Store>) -> Self {
90        Self {
91            store,
92            interval: DEFAULT_TICK_INTERVAL,
93        }
94    }
95
96    /// Set how often due schedules are polled.
97    pub fn interval(mut self, interval: Duration) -> Self {
98        self.interval = interval;
99        self
100    }
101
102    /// Run the tick loop until `shutdown` is cancelled.
103    ///
104    /// Store errors are logged per-schedule; the loop keeps going.
105    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    /// Process one batch of due schedules.
128    ///
129    /// Exposed for tests and for callers that drive the tick themselves.
130    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}