Skip to main content

ironflow_store/memory/
schedule_store.rs

1//! In-memory [`ScheduleStore`] implementation.
2
3use chrono::Utc;
4use uuid::Uuid;
5
6use crate::entities::{NewSchedule, Page, Schedule, ScheduleUpdate};
7use crate::error::StoreError;
8use crate::memory::InMemoryStore;
9use crate::schedule_store::ScheduleStore;
10use crate::store::StoreFuture;
11
12impl ScheduleStore for InMemoryStore {
13    fn create_schedule(&self, req: NewSchedule) -> StoreFuture<'_, Schedule> {
14        Box::pin(async move {
15            let now = Utc::now();
16            let schedule = Schedule {
17                id: Uuid::now_v7(),
18                workflow_name: req.workflow_name,
19                cron_expression: req.cron_expression,
20                inputs: req.inputs,
21                source: req.source,
22                disabled_at: None,
23                last_triggered_at: None,
24                next_trigger_at: req.next_trigger_at,
25                created_by_user_id: req.created_by_user_id,
26                created_at: now,
27                updated_at: now,
28            };
29            let mut state = self.state.write().await;
30            state.schedules.insert(schedule.id, schedule.clone());
31            Ok(schedule)
32        })
33    }
34
35    fn find_schedule_by_id(&self, id: Uuid) -> StoreFuture<'_, Option<Schedule>> {
36        Box::pin(async move {
37            let state = self.state.read().await;
38            Ok(state.schedules.get(&id).cloned())
39        })
40    }
41
42    fn list_schedules(&self, page: u32, per_page: u32) -> StoreFuture<'_, Page<Schedule>> {
43        Box::pin(async move {
44            let state = self.state.read().await;
45            let mut all: Vec<_> = state.schedules.values().cloned().collect();
46            all.sort_by_key(|s| std::cmp::Reverse(s.created_at));
47            let total = all.len() as u64;
48            let start = ((page.saturating_sub(1)) as usize) * (per_page as usize);
49            let items: Vec<_> = all
50                .into_iter()
51                .skip(start)
52                .take(per_page as usize)
53                .collect();
54            Ok(Page {
55                items,
56                total,
57                page,
58                per_page,
59            })
60        })
61    }
62
63    fn update_schedule(&self, id: Uuid, update: ScheduleUpdate) -> StoreFuture<'_, Schedule> {
64        Box::pin(async move {
65            let mut state = self.state.write().await;
66            let schedule = state
67                .schedules
68                .get_mut(&id)
69                .ok_or(StoreError::ScheduleNotFound(id))?;
70
71            if let Some(cron) = update.cron_expression {
72                schedule.cron_expression = cron;
73            }
74            if let Some(inputs) = update.inputs {
75                schedule.inputs = inputs;
76            }
77            if let Some(disabled) = update.disabled_at {
78                schedule.disabled_at = disabled;
79            }
80            if let Some(next) = update.next_trigger_at {
81                schedule.next_trigger_at = next;
82            }
83            if let Some(last) = update.last_triggered_at {
84                schedule.last_triggered_at = last;
85            }
86            schedule.updated_at = Utc::now();
87            Ok(schedule.clone())
88        })
89    }
90
91    fn delete_schedule(&self, id: Uuid) -> StoreFuture<'_, ()> {
92        Box::pin(async move {
93            let mut state = self.state.write().await;
94            state
95                .schedules
96                .remove(&id)
97                .ok_or(StoreError::ScheduleNotFound(id))?;
98            Ok(())
99        })
100    }
101
102    fn claim_due_schedules(&self) -> StoreFuture<'_, Vec<Schedule>> {
103        Box::pin(async move {
104            let now = Utc::now();
105            let mut state = self.state.write().await;
106            let due_ids: Vec<Uuid> = state
107                .schedules
108                .values()
109                .filter(|s| s.is_active() && s.next_trigger_at.is_some_and(|at| at <= now))
110                .map(|s| s.id)
111                .collect();
112
113            let mut claimed = Vec::with_capacity(due_ids.len());
114            for id in due_ids {
115                if let Some(schedule) = state.schedules.get_mut(&id) {
116                    schedule.last_triggered_at = Some(now);
117                    schedule.next_trigger_at = None;
118                    schedule.updated_at = now;
119                    claimed.push(schedule.clone());
120                }
121            }
122            Ok(claimed)
123        })
124    }
125}
126
127#[cfg(test)]
128mod tests {
129    use serde_json::json;
130
131    use crate::entities::ScheduleSource;
132
133    use super::*;
134
135    fn new_schedule(workflow: &str, cron: &str) -> NewSchedule {
136        NewSchedule {
137            workflow_name: workflow.to_string(),
138            cron_expression: cron.to_string(),
139            inputs: json!({}),
140            source: ScheduleSource::Api,
141            created_by_user_id: Uuid::now_v7(),
142            next_trigger_at: Some(Utc::now()),
143        }
144    }
145
146    #[tokio::test]
147    async fn create_and_find() {
148        let store = InMemoryStore::new();
149        let created = store
150            .create_schedule(new_schedule("deploy", "0 0 * * * *"))
151            .await
152            .expect("create");
153        assert_eq!(created.workflow_name, "deploy");
154        assert!(created.is_active());
155        assert_eq!(created.source, ScheduleSource::Api);
156
157        let found = store
158            .find_schedule_by_id(created.id)
159            .await
160            .expect("find")
161            .expect("some");
162        assert_eq!(found.id, created.id);
163    }
164
165    #[tokio::test]
166    async fn create_handler_source() {
167        let store = InMemoryStore::new();
168        let created = store
169            .create_schedule(NewSchedule {
170                source: ScheduleSource::Handler,
171                ..new_schedule("nightly", "0 0 * * *")
172            })
173            .await
174            .expect("create");
175        assert_eq!(created.source, ScheduleSource::Handler);
176    }
177
178    #[tokio::test]
179    async fn list_paginated() {
180        let store = InMemoryStore::new();
181        for i in 0..5 {
182            store
183                .create_schedule(new_schedule(&format!("wf-{i}"), "0 0 * * * *"))
184                .await
185                .expect("create");
186        }
187        let page = store.list_schedules(1, 3).await.expect("list");
188        assert_eq!(page.items.len(), 3);
189        assert_eq!(page.total, 5);
190
191        let page2 = store.list_schedules(2, 3).await.expect("list");
192        assert_eq!(page2.items.len(), 2);
193    }
194
195    #[tokio::test]
196    async fn update_fields() {
197        let store = InMemoryStore::new();
198        let created = store
199            .create_schedule(new_schedule("deploy", "0 0 * * * *"))
200            .await
201            .expect("create");
202
203        let updated = store
204            .update_schedule(
205                created.id,
206                ScheduleUpdate {
207                    disabled_at: Some(Some(Utc::now())),
208                    cron_expression: Some("0 30 * * * *".to_string()),
209                    ..Default::default()
210                },
211            )
212            .await
213            .expect("update");
214
215        assert!(!updated.is_active());
216        assert_eq!(updated.cron_expression, "0 30 * * * *");
217    }
218
219    #[tokio::test]
220    async fn update_not_found() {
221        let store = InMemoryStore::new();
222        let err = store
223            .update_schedule(Uuid::now_v7(), ScheduleUpdate::default())
224            .await
225            .unwrap_err();
226        assert!(matches!(err, StoreError::ScheduleNotFound(_)));
227    }
228
229    #[tokio::test]
230    async fn delete_existing() {
231        let store = InMemoryStore::new();
232        let created = store
233            .create_schedule(new_schedule("deploy", "0 0 * * * *"))
234            .await
235            .expect("create");
236
237        store.delete_schedule(created.id).await.expect("delete");
238
239        let found = store.find_schedule_by_id(created.id).await.expect("find");
240        assert!(found.is_none());
241    }
242
243    #[tokio::test]
244    async fn delete_not_found() {
245        let store = InMemoryStore::new();
246        let err = store.delete_schedule(Uuid::now_v7()).await.unwrap_err();
247        assert!(matches!(err, StoreError::ScheduleNotFound(_)));
248    }
249
250    #[tokio::test]
251    async fn claim_due_schedules_filters_correctly() {
252        use chrono::TimeDelta;
253
254        let store = InMemoryStore::new();
255
256        let past = store
257            .create_schedule(NewSchedule {
258                workflow_name: "past".to_string(),
259                cron_expression: "0 0 * * * *".to_string(),
260                inputs: json!({}),
261                source: ScheduleSource::Api,
262                created_by_user_id: Uuid::now_v7(),
263                next_trigger_at: Some(Utc::now() - TimeDelta::seconds(60)),
264            })
265            .await
266            .expect("create past");
267
268        let _future = store
269            .create_schedule(NewSchedule {
270                workflow_name: "future".to_string(),
271                cron_expression: "0 0 * * * *".to_string(),
272                inputs: json!({}),
273                source: ScheduleSource::Api,
274                created_by_user_id: Uuid::now_v7(),
275                next_trigger_at: Some(Utc::now() + TimeDelta::seconds(3600)),
276            })
277            .await
278            .expect("create future");
279
280        let disabled = store
281            .create_schedule(NewSchedule {
282                workflow_name: "disabled".to_string(),
283                cron_expression: "0 0 * * * *".to_string(),
284                inputs: json!({}),
285                source: ScheduleSource::Api,
286                created_by_user_id: Uuid::now_v7(),
287                next_trigger_at: Some(Utc::now() - TimeDelta::seconds(60)),
288            })
289            .await
290            .expect("create disabled");
291        store
292            .update_schedule(
293                disabled.id,
294                ScheduleUpdate {
295                    disabled_at: Some(Some(Utc::now())),
296                    ..Default::default()
297                },
298            )
299            .await
300            .expect("disable");
301
302        let due = store.claim_due_schedules().await.expect("claim");
303        assert_eq!(due.len(), 1);
304        assert_eq!(due[0].id, past.id);
305        assert!(due[0].last_triggered_at.is_some());
306        assert!(due[0].next_trigger_at.is_none());
307
308        // Claiming again returns nothing (next_trigger_at was cleared).
309        let due_again = store.claim_due_schedules().await.expect("claim again");
310        assert!(due_again.is_empty());
311    }
312}