ironflow_store/memory/
schedule_store.rs1use 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 let due_again = store.claim_due_schedules().await.expect("claim again");
310 assert!(due_again.is_empty());
311 }
312}