1use 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
15pub 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 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 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 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 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 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}