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