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::{
8    NewSchedule, Page, Schedule, ScheduleFiring, ScheduleFiringPlan, ScheduleNext, ScheduleUpdate,
9    ScheduledRun,
10};
11use crate::error::StoreError;
12use crate::memory::InMemoryStore;
13use crate::schedule_store::ScheduleStore;
14use crate::store::StoreFuture;
15
16impl ScheduleStore for InMemoryStore {
17    fn create_schedule(&self, req: NewSchedule) -> StoreFuture<'_, Schedule> {
18        Box::pin(async move {
19            let now = Utc::now();
20            let schedule = Schedule {
21                id: Uuid::now_v7(),
22                workflow_name: req.workflow_name,
23                cron_expression: req.cron_expression,
24                inputs: req.inputs,
25                source: req.source,
26                disabled_at: None,
27                last_triggered_at: None,
28                next_trigger_at: req.next_trigger_at,
29                last_error: None,
30                priority: req.priority,
31                policy: req.policy,
32                created_by_user_id: req.created_by_user_id,
33                created_at: now,
34                updated_at: now,
35            };
36            let mut state = self.state.write().await;
37            state.schedules.insert(schedule.id, schedule.clone());
38            Ok(schedule)
39        })
40    }
41
42    fn find_schedule_by_id(&self, id: Uuid) -> StoreFuture<'_, Option<Schedule>> {
43        Box::pin(async move {
44            let state = self.state.read().await;
45            Ok(state.schedules.get(&id).cloned())
46        })
47    }
48
49    fn list_schedules(&self, page: u32, per_page: u32) -> StoreFuture<'_, Page<Schedule>> {
50        Box::pin(async move {
51            let state = self.state.read().await;
52            let mut all: Vec<_> = state.schedules.values().cloned().collect();
53            all.sort_by_key(|s| std::cmp::Reverse(s.created_at));
54            let total = all.len() as u64;
55            let start = ((page.saturating_sub(1)) as usize) * (per_page as usize);
56            let items: Vec<_> = all
57                .into_iter()
58                .skip(start)
59                .take(per_page as usize)
60                .collect();
61            Ok(Page {
62                items,
63                total,
64                page,
65                per_page,
66            })
67        })
68    }
69
70    fn update_schedule(&self, id: Uuid, update: ScheduleUpdate) -> StoreFuture<'_, Schedule> {
71        Box::pin(async move {
72            let mut state = self.state.write().await;
73            let schedule = state
74                .schedules
75                .get_mut(&id)
76                .ok_or(StoreError::ScheduleNotFound(id))?;
77
78            if let Some(cron) = update.cron_expression {
79                schedule.cron_expression = cron;
80            }
81            if let Some(inputs) = update.inputs {
82                schedule.inputs = inputs;
83            }
84            if let Some(disabled) = update.disabled_at {
85                schedule.disabled_at = disabled;
86            }
87            if let Some(next) = update.next_trigger_at {
88                schedule.next_trigger_at = next;
89            }
90            if let Some(last) = update.last_triggered_at {
91                schedule.last_triggered_at = last;
92            }
93            if let Some(error) = update.last_error {
94                schedule.last_error = error;
95            }
96            if let Some(priority) = update.priority {
97                schedule.priority = priority;
98            }
99            if let Some(policy) = update.policy {
100                schedule.policy = policy;
101            }
102            schedule.updated_at = Utc::now();
103            Ok(schedule.clone())
104        })
105    }
106
107    fn delete_schedule(&self, id: Uuid) -> StoreFuture<'_, ()> {
108        Box::pin(async move {
109            let mut state = self.state.write().await;
110            state
111                .schedules
112                .remove(&id)
113                .ok_or(StoreError::ScheduleNotFound(id))?;
114            Ok(())
115        })
116    }
117
118    fn list_due_schedules(&self) -> StoreFuture<'_, Vec<Schedule>> {
119        Box::pin(async move {
120            let now = Utc::now();
121            let state = self.state.read().await;
122            let mut due: Vec<Schedule> = state
123                .schedules
124                .values()
125                .filter(|s| s.is_active() && s.next_trigger_at.is_some_and(|at| at <= now))
126                .cloned()
127                .collect();
128            due.sort_by_key(|s| s.next_trigger_at);
129            Ok(due)
130        })
131    }
132
133    fn fire_due_schedule(
134        &self,
135        id: Uuid,
136        due: DateTime<Utc>,
137        plan: ScheduleFiringPlan,
138    ) -> StoreFuture<'_, Option<ScheduleFiring>> {
139        Box::pin(async move {
140            // One write lock covers the check, the runs and the schedule update.
141            let mut state = self.state.write().await;
142            let Some(schedule) = state
143                .schedules
144                .get(&id)
145                .filter(|s| s.is_active() && s.next_trigger_at == Some(due))
146                .cloned()
147            else {
148                return Ok(None);
149            };
150
151            // Every run is built from the same schedule, so an error other
152            // than a concurrency conflict (an invalid priority, say) fails on
153            // the first occurrence, before anything is written.
154            let mut runs = Vec::with_capacity(plan.occurrences.len());
155            let mut overlapped = Vec::new();
156            for occurrence in plan.occurrences {
157                let mut new_run = schedule.new_run(Some(occurrence), None);
158                new_run.idempotency_key = Some(Schedule::occurrence_key(id, occurrence));
159                match insert_run(&mut state, new_run) {
160                    Ok(run) => runs.push(ScheduledRun { occurrence, run }),
161                    Err(StoreError::ConcurrencyConflict { .. }) => overlapped.push(occurrence),
162                    Err(e) => return Err(e),
163                }
164            }
165
166            let now = Utc::now();
167            let schedule = state
168                .schedules
169                .get_mut(&id)
170                .ok_or(StoreError::ScheduleNotFound(id))?;
171            if !runs.is_empty() {
172                schedule.last_triggered_at = Some(now);
173            }
174            schedule.updated_at = now;
175            match plan.next {
176                ScheduleNext::At(at) => schedule.next_trigger_at = Some(at),
177                ScheduleNext::Disable { error } => {
178                    schedule.next_trigger_at = None;
179                    schedule.disabled_at = Some(now);
180                    schedule.last_error = Some(error);
181                }
182            }
183
184            Ok(Some(ScheduleFiring {
185                schedule: schedule.clone(),
186                runs,
187                overlapped,
188            }))
189        })
190    }
191}
192
193#[cfg(test)]
194mod tests {
195    use chrono::TimeDelta;
196    use chrono_tz::Tz;
197    use serde_json::json;
198
199    use crate::entities::{
200        CatchupPolicy, OverlapPolicy, RunFilter, SchedulePolicy, ScheduleSource, TriggerKind,
201    };
202    use crate::store::RunStore;
203
204    use super::*;
205
206    fn new_schedule(workflow: &str, cron: &str) -> NewSchedule {
207        NewSchedule {
208            workflow_name: workflow.to_string(),
209            cron_expression: cron.to_string(),
210            inputs: json!({}),
211            source: ScheduleSource::Api,
212            priority: 0,
213            policy: SchedulePolicy::default(),
214            created_by_user_id: Some(Uuid::now_v7()),
215            next_trigger_at: Some(Utc::now()),
216        }
217    }
218
219    fn once(occurrence: DateTime<Utc>, next: ScheduleNext) -> ScheduleFiringPlan {
220        ScheduleFiringPlan {
221            occurrences: vec![occurrence],
222            next,
223        }
224    }
225
226    fn in_an_hour() -> ScheduleNext {
227        ScheduleNext::At(Utc::now() + TimeDelta::seconds(3600))
228    }
229
230    #[tokio::test]
231    async fn create_and_find() {
232        let store = InMemoryStore::new();
233        let created = store
234            .create_schedule(new_schedule("deploy", "0 0 * * * *"))
235            .await
236            .expect("create");
237        assert_eq!(created.workflow_name, "deploy");
238        assert!(created.is_active());
239        assert_eq!(created.source, ScheduleSource::Api);
240
241        let found = store
242            .find_schedule_by_id(created.id)
243            .await
244            .expect("find")
245            .expect("some");
246        assert_eq!(found.id, created.id);
247    }
248
249    #[tokio::test]
250    async fn create_handler_source() {
251        let store = InMemoryStore::new();
252        let created = store
253            .create_schedule(NewSchedule {
254                source: ScheduleSource::Handler,
255                ..new_schedule("nightly", "0 0 * * *")
256            })
257            .await
258            .expect("create");
259        assert_eq!(created.source, ScheduleSource::Handler);
260    }
261
262    #[tokio::test]
263    async fn list_paginated() {
264        let store = InMemoryStore::new();
265        for i in 0..5 {
266            store
267                .create_schedule(new_schedule(&format!("wf-{i}"), "0 0 * * * *"))
268                .await
269                .expect("create");
270        }
271        let page = store.list_schedules(1, 3).await.expect("list");
272        assert_eq!(page.items.len(), 3);
273        assert_eq!(page.total, 5);
274
275        let page2 = store.list_schedules(2, 3).await.expect("list");
276        assert_eq!(page2.items.len(), 2);
277    }
278
279    #[tokio::test]
280    async fn update_fields() {
281        let store = InMemoryStore::new();
282        let created = store
283            .create_schedule(new_schedule("deploy", "0 0 * * * *"))
284            .await
285            .expect("create");
286
287        let updated = store
288            .update_schedule(
289                created.id,
290                ScheduleUpdate {
291                    disabled_at: Some(Some(Utc::now())),
292                    cron_expression: Some("0 30 * * * *".to_string()),
293                    ..Default::default()
294                },
295            )
296            .await
297            .expect("update");
298
299        assert!(!updated.is_active());
300        assert_eq!(updated.cron_expression, "0 30 * * * *");
301    }
302
303    #[tokio::test]
304    async fn update_schedule_priority() {
305        let store = InMemoryStore::new();
306        let created = store
307            .create_schedule(NewSchedule {
308                priority: 5,
309                ..new_schedule("deploy", "0 0 * * * *")
310            })
311            .await
312            .expect("create");
313        assert_eq!(created.priority, 5);
314
315        let updated = store
316            .update_schedule(
317                created.id,
318                ScheduleUpdate {
319                    priority: Some(-20),
320                    ..Default::default()
321                },
322            )
323            .await
324            .expect("update");
325        assert_eq!(updated.priority, -20);
326    }
327
328    #[tokio::test]
329    async fn fire_due_schedule_run_inherits_schedule_priority() {
330        let store = InMemoryStore::new();
331        let schedule = store
332            .create_schedule(NewSchedule {
333                priority: 30,
334                ..due_schedule("deploy")
335            })
336            .await
337            .expect("create");
338        let occurrence = schedule.next_trigger_at.expect("due");
339
340        let firing = store
341            .fire_due_schedule(schedule.id, occurrence, once(occurrence, in_an_hour()))
342            .await
343            .expect("fire")
344            .expect("due occurrence fires");
345
346        assert_eq!(firing.runs[0].run.run().priority, 30);
347    }
348
349    #[tokio::test]
350    async fn update_not_found() {
351        let store = InMemoryStore::new();
352        let err = store
353            .update_schedule(Uuid::now_v7(), ScheduleUpdate::default())
354            .await
355            .unwrap_err();
356        assert!(matches!(err, StoreError::ScheduleNotFound(_)));
357    }
358
359    #[tokio::test]
360    async fn delete_existing() {
361        let store = InMemoryStore::new();
362        let created = store
363            .create_schedule(new_schedule("deploy", "0 0 * * * *"))
364            .await
365            .expect("create");
366
367        store.delete_schedule(created.id).await.expect("delete");
368
369        let found = store.find_schedule_by_id(created.id).await.expect("find");
370        assert!(found.is_none());
371    }
372
373    #[tokio::test]
374    async fn delete_not_found() {
375        let store = InMemoryStore::new();
376        let err = store.delete_schedule(Uuid::now_v7()).await.unwrap_err();
377        assert!(matches!(err, StoreError::ScheduleNotFound(_)));
378    }
379
380    fn due_schedule(workflow: &str) -> NewSchedule {
381        NewSchedule {
382            next_trigger_at: Some(Utc::now() - TimeDelta::seconds(60)),
383            ..new_schedule(workflow, "0 0 * * * *")
384        }
385    }
386
387    async fn run_count(store: &InMemoryStore) -> usize {
388        store
389            .list_runs(RunFilter::default(), 1, 100)
390            .await
391            .expect("list runs")
392            .items
393            .len()
394    }
395
396    #[tokio::test]
397    async fn list_due_schedules_filters_and_changes_nothing() {
398        let store = InMemoryStore::new();
399
400        let past = store
401            .create_schedule(due_schedule("past"))
402            .await
403            .expect("create past");
404        store
405            .create_schedule(NewSchedule {
406                next_trigger_at: Some(Utc::now() + TimeDelta::seconds(3600)),
407                ..new_schedule("future", "0 0 * * * *")
408            })
409            .await
410            .expect("create future");
411        let disabled = store
412            .create_schedule(due_schedule("disabled"))
413            .await
414            .expect("create disabled");
415        store
416            .update_schedule(
417                disabled.id,
418                ScheduleUpdate {
419                    disabled_at: Some(Some(Utc::now())),
420                    ..Default::default()
421                },
422            )
423            .await
424            .expect("disable");
425
426        let due = store.list_due_schedules().await.expect("list due");
427        assert_eq!(due.len(), 1);
428        assert_eq!(due[0].id, past.id);
429
430        // Listing is read-only: the schedule is still due at its occurrence.
431        let again = store.list_due_schedules().await.expect("list due again");
432        assert_eq!(again.len(), 1);
433        assert_eq!(again[0].next_trigger_at, past.next_trigger_at);
434        assert!(again[0].last_triggered_at.is_none());
435    }
436
437    #[tokio::test]
438    async fn fire_due_schedule_creates_run_and_advances_next_trigger() {
439        let store = InMemoryStore::new();
440        let schedule = store
441            .create_schedule(NewSchedule {
442                inputs: json!({"env": "prod"}),
443                ..due_schedule("deploy")
444            })
445            .await
446            .expect("create");
447        let occurrence = schedule.next_trigger_at.expect("due");
448        let next = Utc::now() + TimeDelta::seconds(3600);
449
450        let firing = store
451            .fire_due_schedule(
452                schedule.id,
453                occurrence,
454                once(occurrence, ScheduleNext::At(next)),
455            )
456            .await
457            .expect("fire")
458            .expect("due occurrence fires");
459
460        assert_eq!(firing.runs.len(), 1);
461        assert!(firing.overlapped.is_empty());
462        assert!(firing.runs[0].run.is_created());
463        assert_eq!(firing.runs[0].occurrence, occurrence);
464        let run = firing.runs[0].run.run();
465        assert_eq!(run.workflow_name, "deploy");
466        assert_eq!(run.payload, json!({"env": "prod"}));
467        assert_eq!(
468            run.idempotency_key.as_deref(),
469            Some(Schedule::occurrence_key(schedule.id, occurrence).as_str())
470        );
471        assert_eq!(firing.schedule.next_trigger_at, Some(next));
472        assert!(firing.schedule.last_triggered_at.is_some());
473        assert!(firing.schedule.is_active());
474
475        let stored = store
476            .find_schedule_by_id(schedule.id)
477            .await
478            .expect("find")
479            .expect("exists");
480        assert_eq!(stored.next_trigger_at, Some(next));
481        assert!(store.list_due_schedules().await.expect("due").is_empty());
482    }
483
484    #[tokio::test]
485    async fn fire_same_occurrence_twice_creates_one_run() {
486        let store = InMemoryStore::new();
487        let schedule = store
488            .create_schedule(due_schedule("deploy"))
489            .await
490            .expect("create");
491        let occurrence = schedule.next_trigger_at.expect("due");
492        let plan = once(occurrence, in_an_hour());
493
494        let first = store
495            .fire_due_schedule(schedule.id, occurrence, plan.clone())
496            .await
497            .expect("first fire");
498        let second = store
499            .fire_due_schedule(schedule.id, occurrence, plan)
500            .await
501            .expect("second fire");
502
503        assert!(first.is_some());
504        assert!(second.is_none(), "the occurrence was already fired");
505        assert_eq!(run_count(&store).await, 1);
506    }
507
508    #[tokio::test]
509    async fn fire_reuses_the_run_already_bound_to_the_occurrence_key() {
510        let store = InMemoryStore::new();
511        let schedule = store
512            .create_schedule(due_schedule("deploy"))
513            .await
514            .expect("create");
515        let occurrence = schedule.next_trigger_at.expect("due");
516        let mut earlier = schedule.new_run(Some(occurrence), None);
517        earlier.idempotency_key = Some(Schedule::occurrence_key(schedule.id, occurrence));
518        let earlier = store.create_run(earlier).await.expect("earlier run");
519
520        let firing = store
521            .fire_due_schedule(schedule.id, occurrence, once(occurrence, in_an_hour()))
522            .await
523            .expect("fire")
524            .expect("due occurrence fires");
525
526        assert!(!firing.runs[0].run.is_created());
527        assert_eq!(firing.runs[0].run.run().id, earlier.run().id);
528        assert_eq!(run_count(&store).await, 1);
529        assert!(firing.schedule.next_trigger_at.is_some());
530    }
531
532    #[tokio::test]
533    async fn fire_paused_schedule_returns_none() {
534        let store = InMemoryStore::new();
535        let schedule = store
536            .create_schedule(due_schedule("deploy"))
537            .await
538            .expect("create");
539        let occurrence = schedule.next_trigger_at.expect("due");
540        store
541            .update_schedule(
542                schedule.id,
543                ScheduleUpdate {
544                    disabled_at: Some(Some(Utc::now())),
545                    ..Default::default()
546                },
547            )
548            .await
549            .expect("pause");
550
551        let fired = store
552            .fire_due_schedule(schedule.id, occurrence, once(occurrence, in_an_hour()))
553            .await
554            .expect("fire");
555
556        assert!(fired.is_none());
557        assert_eq!(run_count(&store).await, 0);
558    }
559
560    #[tokio::test]
561    async fn fire_unknown_schedule_returns_none() {
562        let store = InMemoryStore::new();
563        let fired = store
564            .fire_due_schedule(Uuid::now_v7(), Utc::now(), once(Utc::now(), in_an_hour()))
565            .await
566            .expect("fire");
567        assert!(fired.is_none());
568    }
569
570    #[tokio::test]
571    async fn fire_with_disable_creates_run_and_disables_with_error() {
572        let store = InMemoryStore::new();
573        let schedule = store
574            .create_schedule(due_schedule("deploy"))
575            .await
576            .expect("create");
577        let occurrence = schedule.next_trigger_at.expect("due");
578
579        let firing = store
580            .fire_due_schedule(
581                schedule.id,
582                occurrence,
583                once(
584                    occurrence,
585                    ScheduleNext::Disable {
586                        error: "no next occurrence".to_string(),
587                    },
588                ),
589            )
590            .await
591            .expect("fire")
592            .expect("due occurrence fires");
593
594        assert!(firing.runs[0].run.is_created());
595        let stored = store
596            .find_schedule_by_id(schedule.id)
597            .await
598            .expect("find")
599            .expect("exists");
600        assert!(!stored.is_active());
601        assert!(stored.next_trigger_at.is_none());
602        assert_eq!(stored.last_error.as_deref(), Some("no next occurrence"));
603    }
604
605    #[tokio::test]
606    async fn update_sets_and_clears_last_error() {
607        let store = InMemoryStore::new();
608        let created = store
609            .create_schedule(new_schedule("deploy", "0 0 * * * *"))
610            .await
611            .expect("create");
612        assert!(created.last_error.is_none());
613
614        let set = store
615            .update_schedule(
616                created.id,
617                ScheduleUpdate {
618                    last_error: Some(Some("boom".to_string())),
619                    ..Default::default()
620                },
621            )
622            .await
623            .expect("set");
624        assert_eq!(set.last_error.as_deref(), Some("boom"));
625
626        let cleared = store
627            .update_schedule(
628                created.id,
629                ScheduleUpdate {
630                    last_error: Some(None),
631                    ..Default::default()
632                },
633            )
634            .await
635            .expect("clear");
636        assert!(cleared.last_error.is_none());
637    }
638
639    #[tokio::test]
640    async fn fire_due_schedule_creates_a_run_per_occurrence() {
641        let store = InMemoryStore::new();
642        let schedule = store
643            .create_schedule(due_schedule("deploy"))
644            .await
645            .expect("create");
646        let due = schedule.next_trigger_at.expect("due");
647        let occurrences = vec![
648            due,
649            due + TimeDelta::seconds(3600),
650            due + TimeDelta::seconds(7200),
651        ];
652
653        let firing = store
654            .fire_due_schedule(
655                schedule.id,
656                due,
657                ScheduleFiringPlan {
658                    occurrences: occurrences.clone(),
659                    next: in_an_hour(),
660                },
661            )
662            .await
663            .expect("fire")
664            .expect("due schedule fires");
665
666        assert!(firing.overlapped.is_empty());
667        let fired: Vec<_> = firing.runs.iter().map(|r| r.occurrence).collect();
668        assert_eq!(fired, occurrences);
669        for scheduled in &firing.runs {
670            let run = scheduled.run.run();
671            assert_eq!(
672                run.trigger,
673                TriggerKind::Cron {
674                    schedule: "0 0 * * * *".to_string(),
675                    schedule_id: Some(schedule.id),
676                    scheduled_for: Some(scheduled.occurrence),
677                }
678            );
679            assert_eq!(
680                run.idempotency_key.as_deref(),
681                Some(Schedule::occurrence_key(schedule.id, scheduled.occurrence).as_str())
682            );
683        }
684        assert_eq!(run_count(&store).await, 3);
685        assert!(firing.schedule.last_triggered_at.is_some());
686    }
687
688    #[tokio::test]
689    async fn fire_due_schedule_with_no_occurrence_only_moves_next_trigger() {
690        let store = InMemoryStore::new();
691        let schedule = store
692            .create_schedule(due_schedule("deploy"))
693            .await
694            .expect("create");
695        let due = schedule.next_trigger_at.expect("due");
696        let next = Utc::now() + TimeDelta::seconds(3600);
697
698        let firing = store
699            .fire_due_schedule(
700                schedule.id,
701                due,
702                ScheduleFiringPlan {
703                    occurrences: Vec::new(),
704                    next: ScheduleNext::At(next),
705                },
706            )
707            .await
708            .expect("fire")
709            .expect("due schedule fires");
710
711        assert!(firing.runs.is_empty());
712        assert_eq!(firing.schedule.next_trigger_at, Some(next));
713        assert!(firing.schedule.last_triggered_at.is_none());
714        assert_eq!(run_count(&store).await, 0);
715    }
716
717    #[tokio::test]
718    async fn fire_due_schedule_reports_overlapped_occurrences() {
719        let store = InMemoryStore::new();
720        let schedule = store
721            .create_schedule(NewSchedule {
722                policy: SchedulePolicy {
723                    overlap: OverlapPolicy::Skip,
724                    ..SchedulePolicy::default()
725                },
726                ..due_schedule("deploy")
727            })
728            .await
729            .expect("create");
730        let due = schedule.next_trigger_at.expect("due");
731        let later = due + TimeDelta::seconds(3600);
732
733        // The first catch-up run holds the schedule key: the second overlaps.
734        let firing = store
735            .fire_due_schedule(
736                schedule.id,
737                due,
738                ScheduleFiringPlan {
739                    occurrences: vec![due, later],
740                    next: ScheduleNext::At(Utc::now() + TimeDelta::seconds(60)),
741                },
742            )
743            .await
744            .expect("fire")
745            .expect("due schedule fires");
746
747        assert_eq!(firing.runs.len(), 1);
748        assert_eq!(firing.runs[0].occurrence, due);
749        assert_eq!(
750            firing.runs[0].run.run().concurrency_key,
751            Some(Schedule::concurrency_key(schedule.id))
752        );
753        assert_eq!(firing.overlapped, vec![later]);
754        assert_eq!(run_count(&store).await, 1);
755    }
756
757    #[tokio::test]
758    async fn fire_with_every_occurrence_overlapped_keeps_last_triggered_at() {
759        let store = InMemoryStore::new();
760        let schedule = store
761            .create_schedule(NewSchedule {
762                policy: SchedulePolicy {
763                    overlap: OverlapPolicy::Skip,
764                    ..SchedulePolicy::default()
765                },
766                ..due_schedule("deploy")
767            })
768            .await
769            .expect("create");
770        let due = schedule.next_trigger_at.expect("due");
771        // A manual trigger of the schedule holds its key.
772        store
773            .create_run(schedule.new_run(None, None))
774            .await
775            .expect("manual run");
776        let next = Utc::now() + TimeDelta::seconds(3600);
777
778        let firing = store
779            .fire_due_schedule(schedule.id, due, once(due, ScheduleNext::At(next)))
780            .await
781            .expect("fire")
782            .expect("due schedule fires");
783
784        assert!(firing.runs.is_empty());
785        assert_eq!(firing.overlapped, vec![due]);
786        assert!(firing.schedule.last_triggered_at.is_none());
787        assert_eq!(firing.schedule.next_trigger_at, Some(next));
788        assert_eq!(run_count(&store).await, 1);
789    }
790
791    #[tokio::test]
792    async fn create_and_update_persist_policy() {
793        let store = InMemoryStore::new();
794        let policy = SchedulePolicy {
795            catchup: CatchupPolicy::All,
796            catchup_max: 3,
797            catchup_window_secs: 7200,
798            overlap: OverlapPolicy::Skip,
799            timezone: Tz::Europe__Paris,
800        };
801        let created = store
802            .create_schedule(NewSchedule {
803                policy: policy.clone(),
804                ..new_schedule("deploy", "0 0 * * * *")
805            })
806            .await
807            .expect("create");
808        assert_eq!(created.policy, policy);
809
810        let changed = SchedulePolicy {
811            catchup: CatchupPolicy::Skip,
812            timezone: Tz::America__New_York,
813            ..policy
814        };
815        let updated = store
816            .update_schedule(
817                created.id,
818                ScheduleUpdate {
819                    policy: Some(changed.clone()),
820                    ..Default::default()
821                },
822            )
823            .await
824            .expect("update");
825        assert_eq!(updated.policy, changed);
826
827        let untouched = store
828            .update_schedule(created.id, ScheduleUpdate::default())
829            .await
830            .expect("no-op update");
831        assert_eq!(untouched.policy, changed);
832    }
833}