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