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