Skip to main content

ironflow_store/memory/
schedule_store.rs

1//! In-memory [`ScheduleStore`] implementation.
2
3use chrono::{DateTime, Utc};
4use uuid::Uuid;
5
6use super::run_store::insert_run;
7use crate::entities::{NewSchedule, Page, Schedule, ScheduleFiring, ScheduleNext, ScheduleUpdate};
8use crate::error::StoreError;
9use crate::memory::InMemoryStore;
10use crate::schedule_store::ScheduleStore;
11use crate::store::StoreFuture;
12
13impl ScheduleStore for InMemoryStore {
14    fn create_schedule(&self, req: NewSchedule) -> StoreFuture<'_, Schedule> {
15        Box::pin(async move {
16            let now = Utc::now();
17            let schedule = Schedule {
18                id: Uuid::now_v7(),
19                workflow_name: req.workflow_name,
20                cron_expression: req.cron_expression,
21                inputs: req.inputs,
22                source: req.source,
23                disabled_at: None,
24                last_triggered_at: None,
25                next_trigger_at: req.next_trigger_at,
26                last_error: None,
27                created_by_user_id: req.created_by_user_id,
28                created_at: now,
29                updated_at: now,
30            };
31            let mut state = self.state.write().await;
32            state.schedules.insert(schedule.id, schedule.clone());
33            Ok(schedule)
34        })
35    }
36
37    fn find_schedule_by_id(&self, id: Uuid) -> StoreFuture<'_, Option<Schedule>> {
38        Box::pin(async move {
39            let state = self.state.read().await;
40            Ok(state.schedules.get(&id).cloned())
41        })
42    }
43
44    fn list_schedules(&self, page: u32, per_page: u32) -> StoreFuture<'_, Page<Schedule>> {
45        Box::pin(async move {
46            let state = self.state.read().await;
47            let mut all: Vec<_> = state.schedules.values().cloned().collect();
48            all.sort_by_key(|s| std::cmp::Reverse(s.created_at));
49            let total = all.len() as u64;
50            let start = ((page.saturating_sub(1)) as usize) * (per_page as usize);
51            let items: Vec<_> = all
52                .into_iter()
53                .skip(start)
54                .take(per_page as usize)
55                .collect();
56            Ok(Page {
57                items,
58                total,
59                page,
60                per_page,
61            })
62        })
63    }
64
65    fn update_schedule(&self, id: Uuid, update: ScheduleUpdate) -> StoreFuture<'_, Schedule> {
66        Box::pin(async move {
67            let mut state = self.state.write().await;
68            let schedule = state
69                .schedules
70                .get_mut(&id)
71                .ok_or(StoreError::ScheduleNotFound(id))?;
72
73            if let Some(cron) = update.cron_expression {
74                schedule.cron_expression = cron;
75            }
76            if let Some(inputs) = update.inputs {
77                schedule.inputs = inputs;
78            }
79            if let Some(disabled) = update.disabled_at {
80                schedule.disabled_at = disabled;
81            }
82            if let Some(next) = update.next_trigger_at {
83                schedule.next_trigger_at = next;
84            }
85            if let Some(last) = update.last_triggered_at {
86                schedule.last_triggered_at = last;
87            }
88            if let Some(error) = update.last_error {
89                schedule.last_error = error;
90            }
91            schedule.updated_at = Utc::now();
92            Ok(schedule.clone())
93        })
94    }
95
96    fn delete_schedule(&self, id: Uuid) -> StoreFuture<'_, ()> {
97        Box::pin(async move {
98            let mut state = self.state.write().await;
99            state
100                .schedules
101                .remove(&id)
102                .ok_or(StoreError::ScheduleNotFound(id))?;
103            Ok(())
104        })
105    }
106
107    fn list_due_schedules(&self) -> StoreFuture<'_, Vec<Schedule>> {
108        Box::pin(async move {
109            let now = Utc::now();
110            let state = self.state.read().await;
111            let mut due: Vec<Schedule> = state
112                .schedules
113                .values()
114                .filter(|s| s.is_active() && s.next_trigger_at.is_some_and(|at| at <= now))
115                .cloned()
116                .collect();
117            due.sort_by_key(|s| s.next_trigger_at);
118            Ok(due)
119        })
120    }
121
122    fn fire_due_schedule(
123        &self,
124        id: Uuid,
125        occurrence: DateTime<Utc>,
126        next: ScheduleNext,
127    ) -> StoreFuture<'_, Option<ScheduleFiring>> {
128        Box::pin(async move {
129            // One write lock covers the check, the run and the schedule update.
130            let mut state = self.state.write().await;
131            let Some(schedule) = state
132                .schedules
133                .get(&id)
134                .filter(|s| s.is_active() && s.next_trigger_at == Some(occurrence))
135            else {
136                return Ok(None);
137            };
138
139            let mut new_run = schedule.new_run(None);
140            new_run.idempotency_key = Some(Schedule::occurrence_key(id, occurrence));
141            let run = insert_run(&mut state, new_run)?;
142
143            let now = Utc::now();
144            let schedule = state
145                .schedules
146                .get_mut(&id)
147                .ok_or(StoreError::ScheduleNotFound(id))?;
148            schedule.last_triggered_at = Some(now);
149            schedule.updated_at = now;
150            match next {
151                ScheduleNext::At(at) => schedule.next_trigger_at = Some(at),
152                ScheduleNext::Disable { error } => {
153                    schedule.next_trigger_at = None;
154                    schedule.disabled_at = Some(now);
155                    schedule.last_error = Some(error);
156                }
157            }
158
159            Ok(Some(ScheduleFiring {
160                schedule: schedule.clone(),
161                run,
162            }))
163        })
164    }
165}
166
167#[cfg(test)]
168mod tests {
169    use chrono::TimeDelta;
170    use serde_json::json;
171
172    use crate::entities::{RunFilter, ScheduleSource};
173    use crate::store::RunStore;
174
175    use super::*;
176
177    fn new_schedule(workflow: &str, cron: &str) -> NewSchedule {
178        NewSchedule {
179            workflow_name: workflow.to_string(),
180            cron_expression: cron.to_string(),
181            inputs: json!({}),
182            source: ScheduleSource::Api,
183            created_by_user_id: Some(Uuid::now_v7()),
184            next_trigger_at: Some(Utc::now()),
185        }
186    }
187
188    #[tokio::test]
189    async fn create_and_find() {
190        let store = InMemoryStore::new();
191        let created = store
192            .create_schedule(new_schedule("deploy", "0 0 * * * *"))
193            .await
194            .expect("create");
195        assert_eq!(created.workflow_name, "deploy");
196        assert!(created.is_active());
197        assert_eq!(created.source, ScheduleSource::Api);
198
199        let found = store
200            .find_schedule_by_id(created.id)
201            .await
202            .expect("find")
203            .expect("some");
204        assert_eq!(found.id, created.id);
205    }
206
207    #[tokio::test]
208    async fn create_handler_source() {
209        let store = InMemoryStore::new();
210        let created = store
211            .create_schedule(NewSchedule {
212                source: ScheduleSource::Handler,
213                ..new_schedule("nightly", "0 0 * * *")
214            })
215            .await
216            .expect("create");
217        assert_eq!(created.source, ScheduleSource::Handler);
218    }
219
220    #[tokio::test]
221    async fn list_paginated() {
222        let store = InMemoryStore::new();
223        for i in 0..5 {
224            store
225                .create_schedule(new_schedule(&format!("wf-{i}"), "0 0 * * * *"))
226                .await
227                .expect("create");
228        }
229        let page = store.list_schedules(1, 3).await.expect("list");
230        assert_eq!(page.items.len(), 3);
231        assert_eq!(page.total, 5);
232
233        let page2 = store.list_schedules(2, 3).await.expect("list");
234        assert_eq!(page2.items.len(), 2);
235    }
236
237    #[tokio::test]
238    async fn update_fields() {
239        let store = InMemoryStore::new();
240        let created = store
241            .create_schedule(new_schedule("deploy", "0 0 * * * *"))
242            .await
243            .expect("create");
244
245        let updated = store
246            .update_schedule(
247                created.id,
248                ScheduleUpdate {
249                    disabled_at: Some(Some(Utc::now())),
250                    cron_expression: Some("0 30 * * * *".to_string()),
251                    ..Default::default()
252                },
253            )
254            .await
255            .expect("update");
256
257        assert!(!updated.is_active());
258        assert_eq!(updated.cron_expression, "0 30 * * * *");
259    }
260
261    #[tokio::test]
262    async fn update_not_found() {
263        let store = InMemoryStore::new();
264        let err = store
265            .update_schedule(Uuid::now_v7(), ScheduleUpdate::default())
266            .await
267            .unwrap_err();
268        assert!(matches!(err, StoreError::ScheduleNotFound(_)));
269    }
270
271    #[tokio::test]
272    async fn delete_existing() {
273        let store = InMemoryStore::new();
274        let created = store
275            .create_schedule(new_schedule("deploy", "0 0 * * * *"))
276            .await
277            .expect("create");
278
279        store.delete_schedule(created.id).await.expect("delete");
280
281        let found = store.find_schedule_by_id(created.id).await.expect("find");
282        assert!(found.is_none());
283    }
284
285    #[tokio::test]
286    async fn delete_not_found() {
287        let store = InMemoryStore::new();
288        let err = store.delete_schedule(Uuid::now_v7()).await.unwrap_err();
289        assert!(matches!(err, StoreError::ScheduleNotFound(_)));
290    }
291
292    fn due_schedule(workflow: &str) -> NewSchedule {
293        NewSchedule {
294            next_trigger_at: Some(Utc::now() - TimeDelta::seconds(60)),
295            ..new_schedule(workflow, "0 0 * * * *")
296        }
297    }
298
299    async fn run_count(store: &InMemoryStore) -> usize {
300        store
301            .list_runs(RunFilter::default(), 1, 100)
302            .await
303            .expect("list runs")
304            .items
305            .len()
306    }
307
308    #[tokio::test]
309    async fn list_due_schedules_filters_and_changes_nothing() {
310        let store = InMemoryStore::new();
311
312        let past = store
313            .create_schedule(due_schedule("past"))
314            .await
315            .expect("create past");
316        store
317            .create_schedule(NewSchedule {
318                next_trigger_at: Some(Utc::now() + TimeDelta::seconds(3600)),
319                ..new_schedule("future", "0 0 * * * *")
320            })
321            .await
322            .expect("create future");
323        let disabled = store
324            .create_schedule(due_schedule("disabled"))
325            .await
326            .expect("create disabled");
327        store
328            .update_schedule(
329                disabled.id,
330                ScheduleUpdate {
331                    disabled_at: Some(Some(Utc::now())),
332                    ..Default::default()
333                },
334            )
335            .await
336            .expect("disable");
337
338        let due = store.list_due_schedules().await.expect("list due");
339        assert_eq!(due.len(), 1);
340        assert_eq!(due[0].id, past.id);
341
342        // Listing is read-only: the schedule is still due at its occurrence.
343        let again = store.list_due_schedules().await.expect("list due again");
344        assert_eq!(again.len(), 1);
345        assert_eq!(again[0].next_trigger_at, past.next_trigger_at);
346        assert!(again[0].last_triggered_at.is_none());
347    }
348
349    #[tokio::test]
350    async fn fire_due_schedule_creates_run_and_advances_next_trigger() {
351        let store = InMemoryStore::new();
352        let schedule = store
353            .create_schedule(NewSchedule {
354                inputs: json!({"env": "prod"}),
355                ..due_schedule("deploy")
356            })
357            .await
358            .expect("create");
359        let occurrence = schedule.next_trigger_at.expect("due");
360        let next = Utc::now() + TimeDelta::seconds(3600);
361
362        let firing = store
363            .fire_due_schedule(schedule.id, occurrence, ScheduleNext::At(next))
364            .await
365            .expect("fire")
366            .expect("due occurrence fires");
367
368        assert!(firing.run.is_created());
369        let run = firing.run.run();
370        assert_eq!(run.workflow_name, "deploy");
371        assert_eq!(run.payload, json!({"env": "prod"}));
372        assert_eq!(
373            run.idempotency_key.as_deref(),
374            Some(Schedule::occurrence_key(schedule.id, occurrence).as_str())
375        );
376        assert_eq!(firing.schedule.next_trigger_at, Some(next));
377        assert!(firing.schedule.last_triggered_at.is_some());
378        assert!(firing.schedule.is_active());
379
380        let stored = store
381            .find_schedule_by_id(schedule.id)
382            .await
383            .expect("find")
384            .expect("exists");
385        assert_eq!(stored.next_trigger_at, Some(next));
386        assert!(store.list_due_schedules().await.expect("due").is_empty());
387    }
388
389    #[tokio::test]
390    async fn fire_same_occurrence_twice_creates_one_run() {
391        let store = InMemoryStore::new();
392        let schedule = store
393            .create_schedule(due_schedule("deploy"))
394            .await
395            .expect("create");
396        let occurrence = schedule.next_trigger_at.expect("due");
397        let next = ScheduleNext::At(Utc::now() + TimeDelta::seconds(3600));
398
399        let first = store
400            .fire_due_schedule(schedule.id, occurrence, next.clone())
401            .await
402            .expect("first fire");
403        let second = store
404            .fire_due_schedule(schedule.id, occurrence, next)
405            .await
406            .expect("second fire");
407
408        assert!(first.is_some());
409        assert!(second.is_none(), "the occurrence was already fired");
410        assert_eq!(run_count(&store).await, 1);
411    }
412
413    #[tokio::test]
414    async fn fire_reuses_the_run_already_bound_to_the_occurrence_key() {
415        let store = InMemoryStore::new();
416        let schedule = store
417            .create_schedule(due_schedule("deploy"))
418            .await
419            .expect("create");
420        let occurrence = schedule.next_trigger_at.expect("due");
421        let mut earlier = schedule.new_run(None);
422        earlier.idempotency_key = Some(Schedule::occurrence_key(schedule.id, occurrence));
423        let earlier = store.create_run(earlier).await.expect("earlier run");
424
425        let firing = store
426            .fire_due_schedule(
427                schedule.id,
428                occurrence,
429                ScheduleNext::At(Utc::now() + TimeDelta::seconds(3600)),
430            )
431            .await
432            .expect("fire")
433            .expect("due occurrence fires");
434
435        assert!(!firing.run.is_created());
436        assert_eq!(firing.run.run().id, earlier.run().id);
437        assert_eq!(run_count(&store).await, 1);
438        assert!(firing.schedule.next_trigger_at.is_some());
439    }
440
441    #[tokio::test]
442    async fn fire_paused_schedule_returns_none() {
443        let store = InMemoryStore::new();
444        let schedule = store
445            .create_schedule(due_schedule("deploy"))
446            .await
447            .expect("create");
448        let occurrence = schedule.next_trigger_at.expect("due");
449        store
450            .update_schedule(
451                schedule.id,
452                ScheduleUpdate {
453                    disabled_at: Some(Some(Utc::now())),
454                    ..Default::default()
455                },
456            )
457            .await
458            .expect("pause");
459
460        let fired = store
461            .fire_due_schedule(
462                schedule.id,
463                occurrence,
464                ScheduleNext::At(Utc::now() + TimeDelta::seconds(3600)),
465            )
466            .await
467            .expect("fire");
468
469        assert!(fired.is_none());
470        assert_eq!(run_count(&store).await, 0);
471    }
472
473    #[tokio::test]
474    async fn fire_unknown_schedule_returns_none() {
475        let store = InMemoryStore::new();
476        let fired = store
477            .fire_due_schedule(Uuid::now_v7(), Utc::now(), ScheduleNext::At(Utc::now()))
478            .await
479            .expect("fire");
480        assert!(fired.is_none());
481    }
482
483    #[tokio::test]
484    async fn fire_with_disable_creates_run_and_disables_with_error() {
485        let store = InMemoryStore::new();
486        let schedule = store
487            .create_schedule(due_schedule("deploy"))
488            .await
489            .expect("create");
490        let occurrence = schedule.next_trigger_at.expect("due");
491
492        let firing = store
493            .fire_due_schedule(
494                schedule.id,
495                occurrence,
496                ScheduleNext::Disable {
497                    error: "no next occurrence".to_string(),
498                },
499            )
500            .await
501            .expect("fire")
502            .expect("due occurrence fires");
503
504        assert!(firing.run.is_created());
505        let stored = store
506            .find_schedule_by_id(schedule.id)
507            .await
508            .expect("find")
509            .expect("exists");
510        assert!(!stored.is_active());
511        assert!(stored.next_trigger_at.is_none());
512        assert_eq!(stored.last_error.as_deref(), Some("no next occurrence"));
513    }
514
515    #[tokio::test]
516    async fn update_sets_and_clears_last_error() {
517        let store = InMemoryStore::new();
518        let created = store
519            .create_schedule(new_schedule("deploy", "0 0 * * * *"))
520            .await
521            .expect("create");
522        assert!(created.last_error.is_none());
523
524        let set = store
525            .update_schedule(
526                created.id,
527                ScheduleUpdate {
528                    last_error: Some(Some("boom".to_string())),
529                    ..Default::default()
530                },
531            )
532            .await
533            .expect("set");
534        assert_eq!(set.last_error.as_deref(), Some("boom"));
535
536        let cleared = store
537            .update_schedule(
538                created.id,
539                ScheduleUpdate {
540                    last_error: Some(None),
541                    ..Default::default()
542                },
543            )
544            .await
545            .expect("clear");
546        assert!(cleared.last_error.is_none());
547    }
548}