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