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;
12use uuid::Uuid;
13
14use crate::schedule_ticker::next_trigger;
15
16pub 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 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 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 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 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 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}