Skip to main content

ironflow_api/
schedule_sync.rs

1//! Reconciliation of handler-declared schedules with the database.
2//!
3//! Call [`sync_handler_schedules`] once at server startup, before spawning
4//! the [`ScheduleTicker`](super::schedule_ticker::ScheduleTicker).
5
6use std::collections::{HashMap, HashSet};
7
8use ironflow_engine::engine::Engine;
9use ironflow_store::entities::{NewSchedule, ScheduleSource, ScheduleUpdate};
10use ironflow_store::store::Store;
11use tracing::info;
12
13use crate::schedule_ticker::next_trigger;
14
15/// Reconcile DB schedules with handler-declared schedules.
16///
17/// Performs a three-way sync at startup:
18/// 1. **Create** DB rows for handlers that declare a schedule but have no
19///    corresponding `source = handler` row.
20/// 2. **Update** the cron expression when the handler's cron changed.
21/// 3. **Delete** orphan `source = handler` rows whose handler was removed
22///    from the code.
23///
24/// # Errors
25///
26/// Returns the first store error encountered. Changes applied before the
27/// error are kept.
28///
29/// # Examples
30///
31/// ```no_run
32/// use std::sync::Arc;
33/// use ironflow_api::schedule_sync::sync_handler_schedules;
34/// use ironflow_engine::engine::Engine;
35/// use ironflow_core::providers::claude::ClaudeCodeProvider;
36/// use ironflow_store::memory::InMemoryStore;
37/// use ironflow_store::store::Store;
38///
39/// # async fn example() -> Result<(), Box<dyn std::error::Error>> {
40/// let store: Arc<dyn Store> = Arc::new(InMemoryStore::new());
41/// let engine = Engine::new(store.clone(), Arc::new(ClaudeCodeProvider::new()));
42/// sync_handler_schedules(&engine, store.as_ref()).await?;
43/// # Ok(())
44/// # }
45/// ```
46pub async fn sync_handler_schedules(
47    engine: &Engine,
48    store: &dyn Store,
49) -> Result<(), ironflow_store::error::StoreError> {
50    let handlers = engine.scheduled_handlers();
51    let handler_map: HashMap<&str, &str> = handlers
52        .iter()
53        .map(|(name, cron)| (*name, cron.as_str()))
54        .collect();
55    let handler_names: HashSet<&str> = handler_map.keys().copied().collect();
56
57    let existing = store.list_schedules(1, 1000).await?;
58
59    let handler_schedules: Vec<_> = existing
60        .items
61        .iter()
62        .filter(|s| s.source == ScheduleSource::Handler)
63        .collect();
64
65    // 1. Delete orphans: handler-declared schedules whose handler no longer exists.
66    for schedule in &handler_schedules {
67        if !handler_names.contains(schedule.workflow_name.as_str()) {
68            store.delete_schedule(schedule.id).await?;
69            info!(
70                workflow = %schedule.workflow_name,
71                "removed orphan handler schedule (handler no longer exists)"
72            );
73        }
74    }
75
76    // Build lookup of remaining handler schedules by workflow name.
77    let existing_map: HashMap<&str, &ironflow_store::entities::Schedule> = handler_schedules
78        .iter()
79        .filter(|s| handler_names.contains(s.workflow_name.as_str()))
80        .map(|s| (s.workflow_name.as_str(), *s))
81        .collect();
82
83    for (name, cron_str) in &handler_map {
84        match existing_map.get(name) {
85            Some(existing) if existing.cron_expression != *cron_str => {
86                // 2. Cron changed in code: update DB row.
87                let next = next_trigger(cron_str).unwrap_or(None);
88                store
89                    .update_schedule(
90                        existing.id,
91                        ScheduleUpdate {
92                            cron_expression: Some(cron_str.to_string()),
93                            next_trigger_at: Some(next),
94                            ..Default::default()
95                        },
96                    )
97                    .await?;
98                info!(
99                    workflow = %name,
100                    old_cron = %existing.cron_expression,
101                    new_cron = %cron_str,
102                    "updated handler schedule cron"
103                );
104            }
105            Some(_) => {}
106            None => {
107                // 3. Missing: create a new handler schedule.
108                let next = next_trigger(cron_str).unwrap_or(None);
109                store
110                    .create_schedule(NewSchedule {
111                        workflow_name: name.to_string(),
112                        cron_expression: cron_str.to_string(),
113                        inputs: serde_json::json!({}),
114                        source: ScheduleSource::Handler,
115                        created_by_user_id: None,
116                        next_trigger_at: next,
117                    })
118                    .await?;
119                info!(
120                    workflow = %name,
121                    schedule = %cron_str,
122                    "synced handler-declared schedule to DB"
123                );
124            }
125        }
126    }
127
128    Ok(())
129}
130
131#[cfg(test)]
132mod tests {
133    use ironflow_core::providers::claude::ClaudeCodeProvider;
134    use ironflow_engine::context::WorkflowContext;
135    use ironflow_engine::handler::{HandlerFuture, WorkflowHandler};
136    use ironflow_engine::prelude::CronSchedule;
137    use ironflow_store::entities::{NewSchedule, ScheduleSource};
138    use ironflow_store::memory::InMemoryStore;
139    use ironflow_store::store::Store;
140    use serde_json::json;
141    use std::sync::Arc;
142    use uuid::Uuid;
143
144    use super::*;
145
146    struct NamedScheduled {
147        wf_name: &'static str,
148        cron: CronSchedule,
149    }
150    impl WorkflowHandler for NamedScheduled {
151        fn name(&self) -> &str {
152            self.wf_name
153        }
154        fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
155            Box::pin(async { Ok(()) })
156        }
157        fn schedule(&self) -> Option<&CronSchedule> {
158            Some(&self.cron)
159        }
160    }
161
162    #[tokio::test]
163    async fn creates_missing_handler_schedules() {
164        let store: Arc<dyn Store> = Arc::new(InMemoryStore::new());
165        let provider = Arc::new(ClaudeCodeProvider::new());
166        let mut engine = ironflow_engine::engine::Engine::new(store.clone(), provider);
167        engine
168            .register(NamedScheduled {
169                wf_name: "nightly-cleanup",
170                cron: CronSchedule::new("0 0 * * *").unwrap(),
171            })
172            .expect("register");
173
174        sync_handler_schedules(&engine, store.as_ref())
175            .await
176            .expect("sync");
177
178        let page = store.list_schedules(1, 10).await.expect("list");
179        assert_eq!(page.items.len(), 1);
180        assert_eq!(page.items[0].workflow_name, "nightly-cleanup");
181        assert_eq!(page.items[0].source, ScheduleSource::Handler);
182    }
183
184    #[tokio::test]
185    async fn skips_existing_handler_schedule() {
186        let store: Arc<dyn Store> = Arc::new(InMemoryStore::new());
187        store
188            .create_schedule(NewSchedule {
189                workflow_name: "deploy".to_string(),
190                cron_expression: "*/5 * * * *".to_string(),
191                inputs: json!({"env": "prod"}),
192                source: ScheduleSource::Handler,
193                created_by_user_id: None,
194                next_trigger_at: None,
195            })
196            .await
197            .expect("seed");
198
199        let provider = Arc::new(ClaudeCodeProvider::new());
200        let mut engine = ironflow_engine::engine::Engine::new(store.clone(), provider);
201        engine
202            .register(NamedScheduled {
203                wf_name: "deploy",
204                cron: CronSchedule::new("*/5 * * * *").unwrap(),
205            })
206            .expect("register");
207
208        sync_handler_schedules(&engine, store.as_ref())
209            .await
210            .expect("sync");
211
212        let page = store.list_schedules(1, 10).await.expect("list");
213        assert_eq!(page.items.len(), 1);
214        assert_eq!(page.items[0].cron_expression, "*/5 * * * *");
215    }
216
217    #[tokio::test]
218    async fn removes_orphan_handler_schedules() {
219        let store: Arc<dyn Store> = Arc::new(InMemoryStore::new());
220        store
221            .create_schedule(NewSchedule {
222                workflow_name: "removed-workflow".to_string(),
223                cron_expression: "0 0 * * *".to_string(),
224                inputs: json!({}),
225                source: ScheduleSource::Handler,
226                created_by_user_id: None,
227                next_trigger_at: None,
228            })
229            .await
230            .expect("seed orphan");
231
232        // API schedule must NOT be removed.
233        store
234            .create_schedule(NewSchedule {
235                workflow_name: "user-created".to_string(),
236                cron_expression: "0 12 * * *".to_string(),
237                inputs: json!({}),
238                source: ScheduleSource::Api,
239                created_by_user_id: Some(Uuid::now_v7()),
240                next_trigger_at: None,
241            })
242            .await
243            .expect("seed api");
244
245        let provider = Arc::new(ClaudeCodeProvider::new());
246        let engine = ironflow_engine::engine::Engine::new(store.clone(), provider);
247
248        sync_handler_schedules(&engine, store.as_ref())
249            .await
250            .expect("sync");
251
252        let page = store.list_schedules(1, 10).await.expect("list");
253        assert_eq!(page.items.len(), 1);
254        assert_eq!(page.items[0].workflow_name, "user-created");
255        assert_eq!(page.items[0].source, ScheduleSource::Api);
256    }
257
258    #[tokio::test]
259    async fn updates_changed_cron() {
260        let store: Arc<dyn Store> = Arc::new(InMemoryStore::new());
261        store
262            .create_schedule(NewSchedule {
263                workflow_name: "deploy".to_string(),
264                cron_expression: "0 0 * * *".to_string(),
265                inputs: json!({}),
266                source: ScheduleSource::Handler,
267                created_by_user_id: None,
268                next_trigger_at: None,
269            })
270            .await
271            .expect("seed");
272
273        let provider = Arc::new(ClaudeCodeProvider::new());
274        let mut engine = ironflow_engine::engine::Engine::new(store.clone(), provider);
275        engine
276            .register(NamedScheduled {
277                wf_name: "deploy",
278                cron: CronSchedule::new("0 */6 * * *").unwrap(),
279            })
280            .expect("register");
281
282        sync_handler_schedules(&engine, store.as_ref())
283            .await
284            .expect("sync");
285
286        let page = store.list_schedules(1, 10).await.expect("list");
287        assert_eq!(page.items.len(), 1);
288        assert_eq!(page.items[0].cron_expression, "0 */6 * * *");
289    }
290}