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