1use std::collections::{BTreeMap, HashMap, HashSet};
2
3use chrono::{DateTime, Duration, Utc};
4use rust_decimal::Decimal;
5use uuid::Uuid;
6
7use crate::entities::{
8 ApiKey, ConcurrencyGroupBacklog, IDEMPOTENCY_WINDOW, LeaseRequest, LeaseUpdate, NewRun,
9 NewStep, NewStepDependency, Page, PurgePolicy, PurgeReason, PurgeableRun, ReapedRun, Run,
10 RunActor, RunCreation, RunFilter, RunStats, RunStatus, RunUpdate, StatsHistoryBucket,
11 StatsHistoryFilter, Step, StepApproval, StepDependency, StepStatus, StepUpdate, TriggerKind,
12 User, validate_concurrency_limits,
13};
14use crate::error::StoreError;
15use crate::store::{LEASE_EXPIRED_ERROR, RunStore, StoreFuture};
16
17use super::descendants::active_descendants;
18use super::stats_history::aggregate_history_buckets;
19use super::{InMemoryStore, State};
20
21fn resolve_created_by_label(
26 actor: Option<&RunActor>,
27 users: &HashMap<Uuid, User>,
28 api_keys: &HashMap<Uuid, ApiKey>,
29) -> Option<String> {
30 let actor = actor?;
31 let username = users.get(&actor.user_id()).map(|u| u.username.clone());
32
33 match actor.api_key_id() {
34 None => username,
35 Some(api_key_id) => {
36 let key_name = api_keys.get(&api_key_id).map(|k| k.name.clone())?;
37 Some(match username {
38 Some(username) => format!("{key_name} ({username})"),
39 None => key_name,
40 })
41 }
42 }
43}
44
45fn run_with_label(run: &Run, state: &State) -> Run {
47 let mut run = run.clone();
48 run.created_by_label =
49 resolve_created_by_label(run.created_by.as_ref(), &state.users, &state.api_keys);
50 run
51}
52
53fn clear_lease(run: &mut Run) {
55 run.worker_id = None;
56 run.lease_expires_at = None;
57}
58
59fn running_count(runs: &HashMap<Uuid, Run>, group: &str) -> u64 {
63 runs.values()
64 .filter(|r| {
65 r.status.state == RunStatus::Running
66 && !matches!(r.trigger, TriggerKind::Workflow)
67 && r.concurrency_limits.iter().any(|l| l.group == group)
68 })
69 .count() as u64
70}
71
72fn is_blocked(run: &Run, runs: &HashMap<Uuid, Run>) -> bool {
75 run.concurrency_limits
76 .iter()
77 .any(|l| running_count(runs, &l.group) >= u64::from(l.limit))
78}
79
80fn is_due(run: &Run, now: DateTime<Utc>) -> bool {
82 matches!(run.status.state, RunStatus::Pending | RunStatus::Retrying)
83 && run.scheduled_at.is_none_or(|at| at <= now)
84}
85
86fn run_matches_filter(run: &Run, filter: &RunFilter, steps: &HashMap<Uuid, Step>) -> bool {
87 if let Some(ref wf) = filter.workflow_name
88 && !run
89 .workflow_name
90 .to_lowercase()
91 .contains(&wf.to_lowercase())
92 {
93 return false;
94 }
95 if let Some(ref status) = filter.status
96 && &run.status.state != status
97 {
98 return false;
99 }
100 if let Some(after) = filter.created_after
101 && run.created_at < after
102 {
103 return false;
104 }
105 if let Some(before) = filter.created_before
106 && run.created_at > before
107 {
108 return false;
109 }
110 if let Some(has_steps) = filter.has_steps
111 && matches!(
112 run.status.state,
113 RunStatus::Completed | RunStatus::Cancelled
114 )
115 {
116 let run_has_steps = steps.values().any(|s| s.run_id == run.id);
117 if has_steps != run_has_steps {
118 return false;
119 }
120 }
121 if let Some(ref labels) = filter.labels {
122 for (key, value) in labels {
123 if run.labels.get(key) != Some(value) {
124 return false;
125 }
126 }
127 }
128 if let Some(user_id) = filter.created_by_user_id
129 && run.created_by.as_ref().map(RunActor::user_id) != Some(user_id)
130 {
131 return false;
132 }
133 if let Some(ref group) = filter.concurrency_group
134 && !run.concurrency_limits.iter().any(|l| &l.group == group)
135 {
136 return false;
137 }
138 true
139}
140
141pub(super) fn insert_run(state: &mut State, req: NewRun) -> Result<RunCreation, StoreError> {
147 validate_concurrency_limits(&req.concurrency_limits)?;
148 let now = Utc::now();
149
150 if let Some(ref key) = req.idempotency_key
151 && let Some(existing) = state
152 .idempotency_keys
153 .get(key)
154 .and_then(|id| state.runs.get(id))
155 {
156 if now - existing.created_at < IDEMPOTENCY_WINDOW {
157 return Ok(RunCreation::Existing(run_with_label(existing, state)));
158 }
159 let stale_id = existing.id;
161 state.idempotency_keys.remove(key);
162 if let Some(stale) = state.runs.get_mut(&stale_id) {
163 stale.idempotency_key = None;
164 }
165 }
166
167 if let Some(ref key) = req.concurrency_key
168 && let Some(holder) = state
169 .runs
170 .values()
171 .filter(|r| {
172 r.concurrency_key.as_deref() == Some(key.as_str()) && !r.status.state.is_terminal()
173 })
174 .min_by_key(|r| r.created_at)
175 {
176 return Err(StoreError::ConcurrencyConflict {
177 key: key.clone(),
178 run_id: holder.id,
179 });
180 }
181
182 let run = Run {
183 id: Uuid::now_v7(),
184 workflow_name: req.workflow_name,
185 status: crate::entities::FsmState::new(RunStatus::Pending, Uuid::now_v7()),
186 trigger: req.trigger,
187 payload: req.payload,
188 error: None,
189 retry_count: 0,
190 max_retries: req.max_retries,
191 cost_usd: Decimal::ZERO,
192 duration_ms: 0,
193 created_at: now,
194 updated_at: now,
195 started_at: None,
196 completed_at: None,
197 handler_version: req.handler_version,
198 labels: req.labels,
199 scheduled_at: req.scheduled_at,
200 created_by: req.created_by,
201 created_by_label: None,
202 idempotency_key: req.idempotency_key.clone(),
203 concurrency_key: req.concurrency_key,
204 concurrency_limits: req.concurrency_limits,
205 max_cost_usd: req.max_cost_usd,
206 worker_id: None,
207 lease_expires_at: None,
208 output: None,
209 lease_recoveries: 0,
210 capacity_wait_kind: None,
211 };
212
213 if let Some(key) = req.idempotency_key {
214 state.idempotency_keys.insert(key, run.id);
215 }
216 state.runs.insert(run.id, run.clone());
217 Ok(RunCreation::Created(run_with_label(&run, state)))
218}
219
220impl RunStore for InMemoryStore {
221 fn create_run(&self, req: NewRun) -> StoreFuture<'_, RunCreation> {
222 Box::pin(async move {
223 let mut state = self.state.write().await;
226 insert_run(&mut state, req)
227 })
228 }
229
230 fn find_run_by_idempotency_key(&self, key: &str) -> StoreFuture<'_, Option<Run>> {
231 let key = key.to_string();
232 Box::pin(async move {
233 let now = Utc::now();
234 let state = self.state.read().await;
235 Ok(state
236 .idempotency_keys
237 .get(&key)
238 .and_then(|id| state.runs.get(id))
239 .filter(|run| now - run.created_at < IDEMPOTENCY_WINDOW)
240 .map(|run| run_with_label(run, &state)))
241 })
242 }
243
244 fn get_run(&self, id: Uuid) -> StoreFuture<'_, Option<Run>> {
245 Box::pin(async move {
246 let state = self.state.read().await;
247 Ok(state.runs.get(&id).map(|r| run_with_label(r, &state)))
248 })
249 }
250
251 fn list_runs(&self, filter: RunFilter, page: u32, per_page: u32) -> StoreFuture<'_, Page<Run>> {
252 Box::pin(async move {
253 let state = self.state.read().await;
254
255 let mut runs: Vec<&Run> = state
256 .runs
257 .values()
258 .filter(|r| run_matches_filter(r, &filter, &state.steps))
259 .collect();
260
261 runs.sort_by_key(|r| std::cmp::Reverse(r.created_at));
263
264 let total = runs.len() as u64;
265 let page = page.max(1);
266 let per_page = per_page.clamp(1, 100);
267 let offset = ((page - 1) * per_page) as usize;
268 let items: Vec<Run> = runs
269 .into_iter()
270 .skip(offset)
271 .take(per_page as usize)
272 .map(|r| run_with_label(r, &state))
273 .collect();
274
275 Ok(Page {
276 items,
277 total,
278 page,
279 per_page,
280 })
281 })
282 }
283
284 fn update_run_status(&self, id: Uuid, new_status: RunStatus) -> StoreFuture<'_, ()> {
285 Box::pin(async move {
286 let mut state = self.state.write().await;
287 let run = state.runs.get_mut(&id).ok_or(StoreError::RunNotFound(id))?;
288
289 if !run.status.state.can_transition_to(&new_status) {
290 return Err(StoreError::InvalidTransition {
291 from: run.status.state,
292 to: new_status,
293 });
294 }
295
296 if run.status.state == new_status && new_status.is_terminal() {
297 return Ok(());
298 }
299
300 let now = Utc::now();
301 run.status.state = new_status;
302 run.updated_at = now;
303
304 if new_status == RunStatus::Running && run.started_at.is_none() {
305 run.started_at = Some(now);
306 }
307 if new_status.is_terminal() {
308 run.completed_at = Some(now);
309 }
310 if new_status != RunStatus::Running {
311 clear_lease(run);
312 }
313 if new_status != RunStatus::Sleeping {
314 run.capacity_wait_kind = None;
315 }
316
317 Ok(())
318 })
319 }
320
321 fn update_run(&self, id: Uuid, update: RunUpdate) -> StoreFuture<'_, ()> {
322 Box::pin(async move {
323 let mut state = self.state.write().await;
324 let run = state.runs.get_mut(&id).ok_or(StoreError::RunNotFound(id))?;
325
326 let now = Utc::now();
327
328 if let Some(status) = update.status {
329 if !run.status.state.can_transition_to(&status) {
330 return Err(StoreError::InvalidTransition {
331 from: run.status.state,
332 to: status,
333 });
334 }
335 if !(run.status.state == status && status.is_terminal()) {
336 run.status.state = status;
337 if status == RunStatus::Running && run.started_at.is_none() {
338 run.started_at = Some(now);
339 }
340 if status.is_terminal() {
341 run.completed_at = Some(now);
342 }
343 if status != RunStatus::Running {
344 clear_lease(run);
345 }
346 }
347 run.capacity_wait_kind = if status == RunStatus::Sleeping {
349 update.capacity_wait_kind.clone()
350 } else {
351 None
352 };
353 }
354
355 match update.lease {
357 Some(LeaseUpdate::Set {
358 worker_id,
359 expires_at,
360 }) => {
361 run.worker_id = Some(worker_id);
362 run.lease_expires_at = Some(expires_at);
363 }
364 Some(LeaseUpdate::Release) => clear_lease(run),
365 None => {}
366 }
367
368 if let Some(error) = update.error {
369 run.error = Some(error);
370 }
371 if update.increment_retry {
372 run.retry_count += 1;
373 }
374 if let Some(cost) = update.cost_usd {
375 run.cost_usd = cost;
376 }
377 if let Some(dur) = update.duration_ms {
378 run.duration_ms = dur;
379 }
380 if let Some(started) = update.started_at {
381 run.started_at = Some(started);
382 }
383 if let Some(completed) = update.completed_at {
384 run.completed_at = Some(completed);
385 }
386 if let Some(scheduled) = update.scheduled_at {
387 run.scheduled_at = Some(scheduled);
388 }
389 if let Some(output) = update.output {
390 run.output = Some(output);
391 }
392
393 run.updated_at = now;
394 Ok(())
395 })
396 }
397
398 fn list_active_descendants(&self, run_id: Uuid) -> StoreFuture<'_, Vec<Run>> {
399 Box::pin(async move {
400 let state = self.state.read().await;
401 Ok(active_descendants(&state.runs, run_id))
402 })
403 }
404
405 fn pick_next_pending(&self, lease: Option<LeaseRequest>) -> StoreFuture<'_, Option<Run>> {
406 Box::pin(async move {
407 let mut state = self.state.write().await;
408 let now = Utc::now();
409
410 let oldest_id = state
418 .runs
419 .values()
420 .filter(|r| is_due(r, now) && !is_blocked(r, &state.runs))
421 .min_by_key(|r| r.created_at)
422 .map(|r| r.id);
423
424 let Some(id) = oldest_id else {
425 return Ok(None);
426 };
427
428 let run = state.runs.get_mut(&id).expect("run exists");
430 let now = Utc::now();
431 run.status.state = RunStatus::Running;
432 run.started_at = Some(now);
433 run.updated_at = now;
434 match lease {
435 Some(lease) => {
436 run.lease_expires_at = Some(lease.expires_at(now));
437 run.worker_id = Some(lease.worker_id);
438 }
439 None => clear_lease(run),
440 }
441 let run = run.clone();
442
443 Ok(Some(run_with_label(&run, &state)))
444 })
445 }
446
447 fn count_blocked_runs_by_group(&self) -> StoreFuture<'_, Vec<ConcurrencyGroupBacklog>> {
448 Box::pin(async move {
449 let state = self.state.read().await;
450 let now = Utc::now();
451
452 let mut blocked: BTreeMap<String, u64> = BTreeMap::new();
453 for run in state.runs.values().filter(|r| is_due(r, now)) {
454 for limit in &run.concurrency_limits {
455 if running_count(&state.runs, &limit.group) >= u64::from(limit.limit) {
456 *blocked.entry(limit.group.clone()).or_default() += 1;
457 }
458 }
459 }
460
461 Ok(blocked
462 .into_iter()
463 .map(|(group, blocked_runs)| ConcurrencyGroupBacklog {
464 group,
465 blocked_runs,
466 })
467 .collect())
468 })
469 }
470
471 fn renew_lease(&self, id: Uuid, lease: LeaseRequest) -> StoreFuture<'_, DateTime<Utc>> {
472 Box::pin(async move {
473 let mut state = self.state.write().await;
474 let run = state.runs.get_mut(&id).ok_or(StoreError::RunNotFound(id))?;
475
476 if run.status.state != RunStatus::Running
477 || run.worker_id.as_deref() != Some(lease.worker_id.as_str())
478 {
479 return Err(StoreError::LeaseLost {
480 run_id: id,
481 held_by: run.worker_id.clone(),
482 });
483 }
484
485 let now = Utc::now();
486 let expires_at = lease.expires_at(now);
487 run.lease_expires_at = Some(expires_at);
488 run.updated_at = now;
489
490 Ok(expires_at)
491 })
492 }
493
494 fn reap_expired_leases(&self, limit: u32) -> StoreFuture<'_, Vec<ReapedRun>> {
495 Box::pin(async move {
496 let mut state = self.state.write().await;
497 let now = Utc::now();
498
499 let mut expired: Vec<Uuid> = state
500 .runs
501 .values()
502 .filter(|r| {
503 r.status.state == RunStatus::Running
504 && r.lease_expires_at.is_some_and(|at| at < now)
505 })
506 .map(|r| r.id)
507 .collect();
508 expired.sort_unstable();
509 expired.truncate(limit as usize);
510
511 let mut reaped = Vec::with_capacity(expired.len());
512 for id in expired {
513 let run = state.runs.get_mut(&id).expect("run exists");
514 run.lease_recoveries += 1;
515 clear_lease(run);
516 run.updated_at = now;
517
518 let to = if run.lease_recoveries > run.max_retries {
519 run.status.state = RunStatus::Failed;
520 run.error = Some(LEASE_EXPIRED_ERROR.to_string());
521 run.completed_at = Some(now);
522 RunStatus::Failed
523 } else {
524 run.status.state = RunStatus::Pending;
525 RunStatus::Pending
526 };
527
528 reaped.push(ReapedRun {
529 run: run.clone(),
530 from: RunStatus::Running,
531 to,
532 });
533 }
534
535 Ok(reaped)
536 })
537 }
538
539 fn claim_due_approval_deadlines(&self, limit: u32) -> StoreFuture<'_, Vec<Step>> {
540 Box::pin(async move {
541 let mut state = self.state.write().await;
542 let now = Utc::now();
543
544 let mut due: Vec<(DateTime<Utc>, Uuid)> = state
545 .steps
546 .values()
547 .filter(|s| {
548 s.status.state == StepStatus::AwaitingApproval
549 && s.approval_deadline_at.is_some_and(|at| at <= now)
550 })
551 .map(|s| (s.approval_deadline_at.expect("deadline is set"), s.id))
552 .collect();
553 due.sort_unstable();
554 due.truncate(limit as usize);
555
556 let mut claimed = Vec::with_capacity(due.len());
557 for (_, id) in due {
558 let step = state.steps.get_mut(&id).expect("step exists");
559 claimed.push(step.clone());
562 step.approval_deadline_at = None;
563 step.updated_at = now;
564 }
565
566 Ok(claimed)
567 })
568 }
569
570 fn claim_due_sleeping_runs(&self, limit: u32) -> StoreFuture<'_, Vec<Run>> {
571 Box::pin(async move {
572 let mut state = self.state.write().await;
573 let now = Utc::now();
574
575 let mut due: Vec<(DateTime<Utc>, Uuid)> = state
576 .runs
577 .values()
578 .filter(|r| r.status.state == RunStatus::Sleeping)
579 .filter_map(|r| r.scheduled_at.filter(|at| *at <= now).map(|at| (at, r.id)))
580 .collect();
581 due.sort_unstable();
582 due.truncate(limit as usize);
583
584 let mut woken = Vec::with_capacity(due.len());
585 for (_, id) in due {
586 let run = state.runs.get_mut(&id).expect("run exists");
587 run.status.state = RunStatus::Pending;
588 run.scheduled_at = None;
589 run.capacity_wait_kind = None;
590 run.updated_at = now;
591 let run = run.clone();
592 woken.push(run_with_label(&run, &state));
593 }
594
595 Ok(woken)
596 })
597 }
598
599 fn list_purgeable_runs(
600 &self,
601 policy: &PurgePolicy,
602 batch_size: u32,
603 ) -> StoreFuture<'_, Vec<PurgeableRun>> {
604 let max_age_days = policy.max_age_days;
605 let max_runs_per_workflow = policy.max_runs_per_workflow;
606 Box::pin(async move {
607 let state = self.state.read().await;
608 let cutoff = Utc::now() - Duration::days(i64::from(max_age_days));
609 let mut result: Vec<PurgeableRun> = Vec::new();
610 let mut seen: HashSet<Uuid> = HashSet::new();
611
612 for run in state.runs.values() {
613 if run.status.state.is_terminal() && run.created_at < cutoff {
614 seen.insert(run.id);
615 result.push(PurgeableRun {
616 run_id: run.id,
617 workflow_name: run.workflow_name.clone(),
618 reason: PurgeReason::TooOld,
619 });
620 }
621 }
622
623 let mut by_workflow: HashMap<&str, Vec<&Run>> = HashMap::new();
624 for run in state.runs.values() {
625 if run.status.state.is_terminal() {
626 by_workflow.entry(&run.workflow_name).or_default().push(run);
627 }
628 }
629 for (_, mut runs) in by_workflow {
630 if runs.len() > max_runs_per_workflow as usize {
631 runs.sort_by_key(|r| r.created_at);
632 let excess = runs.len() - max_runs_per_workflow as usize;
633 for run in runs.into_iter().take(excess) {
634 if seen.insert(run.id) {
635 result.push(PurgeableRun {
636 run_id: run.id,
637 workflow_name: run.workflow_name.clone(),
638 reason: PurgeReason::ExceedsWorkflowLimit,
639 });
640 }
641 }
642 }
643 }
644
645 result.sort_by_key(|p| p.run_id);
646 result.truncate(batch_size as usize);
647 Ok(result)
648 })
649 }
650
651 fn delete_run(&self, id: Uuid) -> StoreFuture<'_, Vec<String>> {
652 Box::pin(async move {
653 let mut state = self.state.write().await;
654
655 if !state.runs.contains_key(&id) {
656 return Err(StoreError::RunNotFound(id));
657 }
658
659 let storage_keys: Vec<String> = state
661 .artifacts
662 .values()
663 .filter(|a| a.run_id == id)
664 .map(|a| a.storage_key.clone())
665 .collect();
666
667 state.artifacts.retain(|_, a| a.run_id != id);
669
670 let step_ids: Vec<Uuid> = state
672 .steps
673 .values()
674 .filter(|s| s.run_id == id)
675 .map(|s| s.id)
676 .collect();
677 state
678 .step_dependencies
679 .retain(|d| !step_ids.contains(&d.step_id) && !step_ids.contains(&d.depends_on));
680
681 state.steps.retain(|_, s| s.run_id != id);
683
684 state.idempotency_keys.retain(|_, &mut run_id| run_id != id);
686
687 state.runs.remove(&id);
689
690 Ok(storage_keys)
691 })
692 }
693
694 fn create_step(&self, req: NewStep) -> StoreFuture<'_, Step> {
695 Box::pin(async move {
696 let mut state = self.state.write().await;
697
698 let attempt = state
699 .runs
700 .get(&req.run_id)
701 .ok_or(StoreError::RunNotFound(req.run_id))?
702 .retry_count
703 + 1;
704
705 let now = Utc::now();
706 let step = Step {
707 id: Uuid::now_v7(),
708 trace_id: req.trace_id,
709 run_id: req.run_id,
710 name: req.name,
711 kind: req.kind,
712 position: req.position,
713 status: crate::entities::FsmState::new(StepStatus::Pending, Uuid::now_v7()),
714 attempt,
715 input: req.input,
716 output: None,
717 error: None,
718 duration_ms: 0,
719 cost_usd: Decimal::ZERO,
720 input_tokens: None,
721 cache_read_input_tokens: None,
722 cache_creation_input_tokens: None,
723 output_tokens: None,
724 created_at: now,
725 updated_at: now,
726 started_at: None,
727 completed_at: None,
728 debug_messages: None,
729 is_error_handler: req.is_error_handler,
730 approval_deadline_at: None,
731 approval_stage: 0,
732 approval_assignee: None,
733 approval_requirement: None,
734 approvals: Vec::new(),
735 account_id: None,
736 environment_id: None,
737 };
738
739 state.steps.insert(step.id, step.clone());
740 Ok(step)
741 })
742 }
743
744 fn update_step(&self, id: Uuid, update: StepUpdate) -> StoreFuture<'_, ()> {
745 Box::pin(async move {
746 let mut state = self.state.write().await;
747 let step = state
748 .steps
749 .get_mut(&id)
750 .ok_or(StoreError::StepNotFound(id))?;
751
752 let now = Utc::now();
753
754 if let Some(status) = update.status {
755 if !matches!(
756 (step.status.state, status),
757 (StepStatus::Pending, StepStatus::Running)
758 | (StepStatus::Pending, StepStatus::Skipped)
759 | (StepStatus::Running, StepStatus::Completed)
760 | (StepStatus::Running, StepStatus::Failed)
761 | (StepStatus::Running, StepStatus::AwaitingApproval)
762 | (StepStatus::AwaitingApproval, StepStatus::Running)
763 | (StepStatus::AwaitingApproval, StepStatus::Completed)
764 | (StepStatus::AwaitingApproval, StepStatus::Failed)
765 | (StepStatus::AwaitingApproval, StepStatus::Rejected)
766 ) {
767 return Err(StoreError::Database(format!(
768 "invalid step status transition: {:?} -> {:?}",
769 step.status.state, status
770 )));
771 }
772 step.status.state = status;
773 }
774 if let Some(output) = update.output {
775 step.output = Some(output);
776 }
777 if let Some(error) = update.error {
778 step.error = Some(error);
779 }
780 if let Some(dur) = update.duration_ms {
781 step.duration_ms = dur;
782 }
783 if let Some(cost) = update.cost_usd {
784 step.cost_usd = cost;
785 }
786 if let Some(tokens) = update.input_tokens {
787 step.input_tokens = Some(tokens);
788 }
789 if let Some(tokens) = update.cache_read_input_tokens {
790 step.cache_read_input_tokens = Some(tokens);
791 }
792 if let Some(tokens) = update.cache_creation_input_tokens {
793 step.cache_creation_input_tokens = Some(tokens);
794 }
795 if let Some(account_id) = update.account_id {
796 step.account_id = Some(account_id);
797 }
798 if let Some(environment_id) = update.environment_id {
799 step.environment_id = Some(environment_id);
800 }
801 if let Some(tokens) = update.output_tokens {
802 step.output_tokens = Some(tokens);
803 }
804 if let Some(started) = update.started_at {
805 step.started_at = Some(started);
806 }
807 if let Some(completed) = update.completed_at {
808 step.completed_at = Some(completed);
809 }
810 if let Some(debug_msgs) = update.debug_messages {
811 step.debug_messages = Some(debug_msgs);
812 }
813 if update.clear_approval_deadline {
816 step.approval_deadline_at = None;
817 } else if let Some(deadline) = update.approval_deadline_at {
818 step.approval_deadline_at = Some(deadline);
819 }
820 if let Some(stage) = update.approval_stage {
821 step.approval_stage = stage;
822 }
823 if let Some(assignee) = update.approval_assignee {
824 step.approval_assignee = Some(assignee);
825 }
826 if let Some(requirement) = update.approval_requirement {
827 step.approval_requirement = Some(requirement);
828 }
829
830 step.updated_at = now;
831 Ok(())
832 })
833 }
834
835 fn get_step(&self, id: Uuid) -> StoreFuture<'_, Option<Step>> {
836 Box::pin(async move {
837 let state = self.state.read().await;
838 Ok(state.steps.get(&id).cloned())
839 })
840 }
841
842 fn record_step_approval(&self, step_id: Uuid, approval: StepApproval) -> StoreFuture<'_, Step> {
843 Box::pin(async move {
844 let mut state = self.state.write().await;
845 let step = state
846 .steps
847 .get_mut(&step_id)
848 .ok_or(StoreError::StepNotFound(step_id))?;
849
850 if !step.approvals.iter().any(|a| a.user_id == approval.user_id) {
851 step.approvals.push(approval);
852 step.updated_at = Utc::now();
853 }
854 Ok(step.clone())
855 })
856 }
857
858 fn list_steps(&self, run_id: Uuid) -> StoreFuture<'_, Vec<Step>> {
859 Box::pin(async move {
860 let state = self.state.read().await;
861 let mut steps: Vec<Step> = state
862 .steps
863 .values()
864 .filter(|s| s.run_id == run_id)
865 .cloned()
866 .collect();
867 steps.sort_by_key(|s| s.position);
868 Ok(steps)
869 })
870 }
871
872 fn get_stats(&self, filter: RunFilter) -> StoreFuture<'_, RunStats> {
873 Box::pin(async move {
874 let state = self.state.read().await;
875
876 let mut total_cost_usd = Decimal::ZERO;
877 let mut total_duration_ms = 0u64;
878 let mut total_runs = 0u64;
879 let mut completed_runs = 0u64;
880 let mut failed_runs = 0u64;
881 let mut cancelled_runs = 0u64;
882 let mut active_runs = 0u64;
883 let mut awaiting_approval_runs = 0u64;
884
885 for run in state.runs.values() {
886 if !run_matches_filter(run, &filter, &state.steps) {
887 continue;
888 }
889
890 total_cost_usd += run.cost_usd;
891 total_duration_ms += run.duration_ms;
892 total_runs += 1;
893
894 match run.status.state {
895 RunStatus::Completed | RunStatus::Warning => completed_runs += 1,
896 RunStatus::Failed => failed_runs += 1,
897 RunStatus::Cancelled => cancelled_runs += 1,
898 RunStatus::AwaitingApproval => {
899 active_runs += 1;
900 awaiting_approval_runs += 1;
901 }
902 RunStatus::Pending
903 | RunStatus::Running
904 | RunStatus::Retrying
905 | RunStatus::Sleeping => {
906 active_runs += 1;
907 }
908 }
909 }
910
911 Ok(RunStats {
912 total_runs,
913 completed_runs,
914 failed_runs,
915 cancelled_runs,
916 active_runs,
917 awaiting_approval_runs,
918 total_cost_usd,
919 total_duration_ms,
920 })
921 })
922 }
923
924 fn get_stats_history(
925 &self,
926 filter: StatsHistoryFilter,
927 ) -> StoreFuture<'_, Vec<StatsHistoryBucket>> {
928 Box::pin(async move {
929 let state = self.state.read().await;
930 let now = Utc::now();
931 let start = now - Duration::hours(filter.period.hours());
932 let run_filter = filter.to_run_filter();
933 let buckets = aggregate_history_buckets(
934 state
935 .runs
936 .values()
937 .filter(|r| run_matches_filter(r, &run_filter, &state.steps)),
938 start,
939 now,
940 filter.granularity,
941 );
942 Ok(buckets)
943 })
944 }
945
946 fn create_step_dependencies(&self, deps: Vec<NewStepDependency>) -> StoreFuture<'_, ()> {
947 Box::pin(async move {
948 let mut state = self.state.write().await;
949
950 for dep in deps {
951 if !state.steps.contains_key(&dep.step_id) {
952 return Err(StoreError::StepNotFound(dep.step_id));
953 }
954 if !state.steps.contains_key(&dep.depends_on) {
955 return Err(StoreError::StepNotFound(dep.depends_on));
956 }
957
958 let already_exists = state
959 .step_dependencies
960 .iter()
961 .any(|d| d.step_id == dep.step_id && d.depends_on == dep.depends_on);
962
963 if !already_exists {
964 state.step_dependencies.push(StepDependency {
965 step_id: dep.step_id,
966 depends_on: dep.depends_on,
967 created_at: Utc::now(),
968 });
969 }
970 }
971
972 Ok(())
973 })
974 }
975
976 fn list_step_dependencies(&self, run_id: Uuid) -> StoreFuture<'_, Vec<StepDependency>> {
977 Box::pin(async move {
978 let state = self.state.read().await;
979
980 let run_step_ids: std::collections::HashSet<Uuid> = state
981 .steps
982 .values()
983 .filter(|s| s.run_id == run_id)
984 .map(|s| s.id)
985 .collect();
986
987 let mut deps: Vec<StepDependency> = state
988 .step_dependencies
989 .iter()
990 .filter(|d| run_step_ids.contains(&d.step_id))
991 .cloned()
992 .collect();
993
994 deps.sort_by_key(|d| d.created_at);
995 Ok(deps)
996 })
997 }
998}
999
1000#[cfg(test)]
1001mod tests {
1002 use std::collections::HashMap;
1003 use std::time::Duration;
1004
1005 use chrono::TimeDelta;
1006 use serde_json::json;
1007 use tokio::spawn;
1008 use tokio::time::sleep;
1009
1010 use super::*;
1011 use crate::api_key_store::ApiKeyStore;
1012 use crate::entities::{
1013 ApiKeyScope, ApiKeyUpdate, ApprovalRequirement, NewApiKey, NewUser, TriggerKind,
1014 };
1015 use crate::user_store::UserStore;
1016
1017 use crate::memory::tests::{create_terminal_run, new_run_req};
1018 use crate::store::RunStore;
1019
1020 use crate::entities::{StepKind, step_trace_id};
1021
1022 fn new_step_req(run_id: Uuid, name: &str, position: u32) -> NewStep {
1023 NewStep {
1024 run_id,
1025 trace_id: step_trace_id(run_id, name, position),
1026 name: name.to_string(),
1027 kind: StepKind::Shell,
1028 position,
1029 input: None,
1030 is_error_handler: false,
1031 }
1032 }
1033
1034 #[tokio::test]
1037 async fn create_run_returns_pending_status() {
1038 let store = InMemoryStore::new();
1039 let run = store
1040 .create_run(new_run_req("test"))
1041 .await
1042 .unwrap()
1043 .into_run();
1044 assert_eq!(run.status.state, RunStatus::Pending);
1045 assert_eq!(run.workflow_name, "test");
1046 assert_eq!(run.retry_count, 0);
1047 assert_eq!(run.max_retries, 3);
1048 }
1049
1050 #[tokio::test]
1051 async fn create_run_generates_unique_ids() {
1052 let store = InMemoryStore::new();
1053 let r1 = store.create_run(new_run_req("a")).await.unwrap().into_run();
1054 let r2 = store.create_run(new_run_req("b")).await.unwrap().into_run();
1055 assert_ne!(r1.id, r2.id);
1056 }
1057
1058 #[tokio::test]
1061 async fn get_run_returns_created_run() {
1062 let store = InMemoryStore::new();
1063 let run = store
1064 .create_run(new_run_req("test"))
1065 .await
1066 .unwrap()
1067 .into_run();
1068 let fetched = store.get_run(run.id).await.unwrap();
1069 assert!(fetched.is_some());
1070 assert_eq!(fetched.unwrap().id, run.id);
1071 }
1072
1073 #[tokio::test]
1074 async fn get_run_returns_none_for_missing() {
1075 let store = InMemoryStore::new();
1076 let fetched = store.get_run(Uuid::nil()).await.unwrap();
1077 assert!(fetched.is_none());
1078 }
1079
1080 #[tokio::test]
1083 async fn update_run_status_valid_transition() {
1084 let store = InMemoryStore::new();
1085 let run = store
1086 .create_run(new_run_req("test"))
1087 .await
1088 .unwrap()
1089 .into_run();
1090
1091 store
1092 .update_run_status(run.id, RunStatus::Running)
1093 .await
1094 .unwrap();
1095
1096 let fetched = store.get_run(run.id).await.unwrap().unwrap();
1097 assert_eq!(fetched.status.state, RunStatus::Running);
1098 assert!(fetched.started_at.is_some());
1099 }
1100
1101 #[tokio::test]
1102 async fn update_run_status_invalid_transition_returns_error() {
1103 let store = InMemoryStore::new();
1104 let run = store
1105 .create_run(new_run_req("test"))
1106 .await
1107 .unwrap()
1108 .into_run();
1109
1110 let result = store.update_run_status(run.id, RunStatus::Completed).await;
1111 assert!(result.is_err());
1112
1113 let err = result.unwrap_err();
1114 assert!(matches!(err, StoreError::InvalidTransition { .. }));
1115 }
1116
1117 #[tokio::test]
1118 async fn update_run_status_not_found() {
1119 let store = InMemoryStore::new();
1120 let result = store
1121 .update_run_status(Uuid::nil(), RunStatus::Running)
1122 .await;
1123 assert!(matches!(result.unwrap_err(), StoreError::RunNotFound(_)));
1124 }
1125
1126 #[tokio::test]
1127 async fn update_run_status_terminal_sets_completed_at() {
1128 let store = InMemoryStore::new();
1129 let run = store
1130 .create_run(new_run_req("test"))
1131 .await
1132 .unwrap()
1133 .into_run();
1134
1135 store
1136 .update_run_status(run.id, RunStatus::Running)
1137 .await
1138 .unwrap();
1139 store
1140 .update_run_status(run.id, RunStatus::Completed)
1141 .await
1142 .unwrap();
1143
1144 let fetched = store.get_run(run.id).await.unwrap().unwrap();
1145 assert_eq!(fetched.status.state, RunStatus::Completed);
1146 assert!(fetched.completed_at.is_some());
1147 }
1148
1149 #[tokio::test]
1150 async fn update_run_status_terminal_to_same_is_idempotent() {
1151 let store = InMemoryStore::new();
1152 let run = store
1153 .create_run(new_run_req("test"))
1154 .await
1155 .unwrap()
1156 .into_run();
1157
1158 store
1159 .update_run_status(run.id, RunStatus::Running)
1160 .await
1161 .unwrap();
1162 store
1163 .update_run_status(run.id, RunStatus::Failed)
1164 .await
1165 .unwrap();
1166
1167 let before = store.get_run(run.id).await.unwrap().unwrap();
1168 let completed_at_before = before.completed_at;
1169
1170 store
1171 .update_run_status(run.id, RunStatus::Failed)
1172 .await
1173 .unwrap();
1174
1175 let after = store.get_run(run.id).await.unwrap().unwrap();
1176 assert_eq!(after.status.state, RunStatus::Failed);
1177 assert_eq!(after.completed_at, completed_at_before);
1178 }
1179
1180 #[tokio::test]
1181 async fn update_run_terminal_to_same_via_update_run_is_idempotent() {
1182 let store = InMemoryStore::new();
1183 let run = store
1184 .create_run(new_run_req("test"))
1185 .await
1186 .unwrap()
1187 .into_run();
1188
1189 store
1190 .update_run_status(run.id, RunStatus::Running)
1191 .await
1192 .unwrap();
1193 store
1194 .update_run(
1195 run.id,
1196 RunUpdate {
1197 status: Some(RunStatus::Failed),
1198 error: Some("first failure".to_string()),
1199 ..RunUpdate::default()
1200 },
1201 )
1202 .await
1203 .unwrap();
1204
1205 let before = store.get_run(run.id).await.unwrap().unwrap();
1206
1207 store
1208 .update_run(
1209 run.id,
1210 RunUpdate {
1211 status: Some(RunStatus::Failed),
1212 ..RunUpdate::default()
1213 },
1214 )
1215 .await
1216 .unwrap();
1217
1218 let after = store.get_run(run.id).await.unwrap().unwrap();
1219 assert_eq!(after.status.state, RunStatus::Failed);
1220 assert_eq!(after.completed_at, before.completed_at);
1221 assert_eq!(after.error, Some("first failure".to_string()));
1222 }
1223
1224 #[tokio::test]
1227 async fn list_runs_empty_store() {
1228 let store = InMemoryStore::new();
1229 let page = store.list_runs(RunFilter::default(), 1, 20).await.unwrap();
1230 assert_eq!(page.total, 0);
1231 assert!(page.items.is_empty());
1232 }
1233
1234 #[tokio::test]
1235 async fn list_runs_with_workflow_filter() {
1236 let store = InMemoryStore::new();
1237 store
1238 .create_run(new_run_req("deploy"))
1239 .await
1240 .unwrap()
1241 .into_run();
1242 store
1243 .create_run(new_run_req("test"))
1244 .await
1245 .unwrap()
1246 .into_run();
1247 store
1248 .create_run(new_run_req("deploy"))
1249 .await
1250 .unwrap()
1251 .into_run();
1252
1253 let filter = RunFilter {
1254 workflow_name: Some("deploy".to_string()),
1255 ..RunFilter::default()
1256 };
1257 let page = store.list_runs(filter, 1, 20).await.unwrap();
1258 assert_eq!(page.total, 2);
1259 assert!(page.items.iter().all(|r| r.workflow_name == "deploy"));
1260 }
1261
1262 #[tokio::test]
1263 async fn list_runs_with_status_filter() {
1264 let store = InMemoryStore::new();
1265 let run = store.create_run(new_run_req("a")).await.unwrap().into_run();
1266 store.create_run(new_run_req("b")).await.unwrap().into_run();
1267
1268 store
1269 .update_run_status(run.id, RunStatus::Running)
1270 .await
1271 .unwrap();
1272
1273 let filter = RunFilter {
1274 status: Some(RunStatus::Running),
1275 ..RunFilter::default()
1276 };
1277 let page = store.list_runs(filter, 1, 20).await.unwrap();
1278 assert_eq!(page.total, 1);
1279 assert_eq!(page.items[0].id, run.id);
1280 }
1281
1282 #[tokio::test]
1283 async fn list_runs_pagination() {
1284 let store = InMemoryStore::new();
1285 for i in 0..5 {
1286 store
1287 .create_run(new_run_req(&format!("wf-{i}")))
1288 .await
1289 .unwrap()
1290 .into_run();
1291 }
1292
1293 let page1 = store.list_runs(RunFilter::default(), 1, 2).await.unwrap();
1294 assert_eq!(page1.total, 5);
1295 assert_eq!(page1.items.len(), 2);
1296 assert_eq!(page1.page, 1);
1297 assert_eq!(page1.per_page, 2);
1298
1299 let page2 = store.list_runs(RunFilter::default(), 2, 2).await.unwrap();
1300 assert_eq!(page2.items.len(), 2);
1301
1302 let page3 = store.list_runs(RunFilter::default(), 3, 2).await.unwrap();
1303 assert_eq!(page3.items.len(), 1);
1304 }
1305
1306 fn lease(worker_id: &str, ttl_secs: u64) -> Option<LeaseRequest> {
1309 Some(LeaseRequest {
1310 worker_id: worker_id.to_string(),
1311 ttl: Duration::from_secs(ttl_secs),
1312 })
1313 }
1314
1315 async fn pick_with_expired_lease(store: &InMemoryStore, max_retries: u32) -> Run {
1320 let mut req = new_run_req("test");
1321 req.max_retries = max_retries;
1322 store.create_run(req).await.unwrap();
1323 let picked = expire_now(store).await;
1324 picked.expect("a pending run was just created")
1325 }
1326
1327 async fn expire_now(store: &InMemoryStore) -> Option<Run> {
1330 let picked = store
1331 .pick_next_pending(Some(LeaseRequest {
1332 worker_id: "worker-1".to_string(),
1333 ttl: Duration::from_nanos(1),
1334 }))
1335 .await
1336 .unwrap();
1337 sleep(Duration::from_millis(2)).await;
1338 picked
1339 }
1340
1341 #[tokio::test]
1342 async fn pick_next_pending_attaches_lease() {
1343 let store = InMemoryStore::new();
1344 store.create_run(new_run_req("test")).await.unwrap();
1345
1346 let picked = store
1347 .pick_next_pending(lease("worker-1", 90))
1348 .await
1349 .unwrap()
1350 .unwrap();
1351
1352 assert_eq!(picked.worker_id.as_deref(), Some("worker-1"));
1353 let expires = picked.lease_expires_at.expect("lease set");
1354 assert!(expires > Utc::now());
1355 assert!(expires <= Utc::now() + TimeDelta::seconds(91));
1356 }
1357
1358 #[tokio::test]
1359 async fn pick_next_pending_without_lease_leaves_run_unowned() {
1360 let store = InMemoryStore::new();
1361 store.create_run(new_run_req("test")).await.unwrap();
1362
1363 let picked = store.pick_next_pending(None).await.unwrap().unwrap();
1364
1365 assert!(picked.worker_id.is_none());
1366 assert!(picked.lease_expires_at.is_none());
1367 }
1368
1369 #[tokio::test]
1370 async fn renew_lease_extends_expiry_for_owner() {
1371 let store = InMemoryStore::new();
1372 store.create_run(new_run_req("test")).await.unwrap();
1373 let picked = store
1374 .pick_next_pending(lease("worker-1", 1))
1375 .await
1376 .unwrap()
1377 .unwrap();
1378
1379 let renewed = store
1380 .renew_lease(picked.id, lease("worker-1", 90).unwrap())
1381 .await
1382 .unwrap();
1383
1384 assert!(renewed > picked.lease_expires_at.unwrap());
1385 let after = store.get_run(picked.id).await.unwrap().unwrap();
1386 assert_eq!(after.lease_expires_at, Some(renewed));
1387 }
1388
1389 #[tokio::test]
1390 async fn renew_lease_rejects_other_worker() {
1391 let store = InMemoryStore::new();
1392 store.create_run(new_run_req("test")).await.unwrap();
1393 let picked = store
1394 .pick_next_pending(lease("worker-1", 90))
1395 .await
1396 .unwrap()
1397 .unwrap();
1398
1399 let err = store
1400 .renew_lease(picked.id, lease("worker-2", 90).unwrap())
1401 .await
1402 .unwrap_err();
1403
1404 assert!(matches!(
1405 err,
1406 StoreError::LeaseLost { held_by: Some(ref w), .. } if w == "worker-1"
1407 ));
1408 }
1409
1410 #[tokio::test]
1411 async fn renew_lease_rejects_run_that_left_running() {
1412 let store = InMemoryStore::new();
1413 store.create_run(new_run_req("test")).await.unwrap();
1414 let picked = store
1415 .pick_next_pending(lease("worker-1", 90))
1416 .await
1417 .unwrap()
1418 .unwrap();
1419 store
1420 .update_run_status(picked.id, RunStatus::Cancelled)
1421 .await
1422 .unwrap();
1423
1424 let err = store
1425 .renew_lease(picked.id, lease("worker-1", 90).unwrap())
1426 .await
1427 .unwrap_err();
1428
1429 assert!(matches!(err, StoreError::LeaseLost { .. }));
1430 }
1431
1432 #[tokio::test]
1433 async fn renew_lease_on_unknown_run_is_not_found() {
1434 let store = InMemoryStore::new();
1435
1436 let err = store
1437 .renew_lease(Uuid::now_v7(), lease("worker-1", 90).unwrap())
1438 .await
1439 .unwrap_err();
1440
1441 assert!(matches!(err, StoreError::RunNotFound(_)));
1442 }
1443
1444 #[tokio::test]
1445 async fn leaving_running_clears_the_lease() {
1446 for target in [
1447 RunStatus::Completed,
1448 RunStatus::Retrying,
1449 RunStatus::AwaitingApproval,
1450 ] {
1451 let store = InMemoryStore::new();
1452 store.create_run(new_run_req("test")).await.unwrap();
1453 let picked = store
1454 .pick_next_pending(lease("worker-1", 90))
1455 .await
1456 .unwrap()
1457 .unwrap();
1458
1459 store.update_run_status(picked.id, target).await.unwrap();
1460
1461 let after = store.get_run(picked.id).await.unwrap().unwrap();
1462 assert!(after.worker_id.is_none(), "worker_id kept for {target}");
1463 assert!(
1464 after.lease_expires_at.is_none(),
1465 "lease_expires_at kept for {target}"
1466 );
1467 }
1468 }
1469
1470 #[tokio::test]
1471 async fn update_run_to_terminal_clears_the_lease() {
1472 let store = InMemoryStore::new();
1473 store.create_run(new_run_req("test")).await.unwrap();
1474 let picked = store
1475 .pick_next_pending(lease("worker-1", 90))
1476 .await
1477 .unwrap()
1478 .unwrap();
1479
1480 store
1481 .update_run(
1482 picked.id,
1483 RunUpdate {
1484 status: Some(RunStatus::Failed),
1485 ..RunUpdate::default()
1486 },
1487 )
1488 .await
1489 .unwrap();
1490
1491 let after = store.get_run(picked.id).await.unwrap().unwrap();
1492 assert!(after.worker_id.is_none());
1493 assert!(after.lease_expires_at.is_none());
1494 }
1495
1496 #[tokio::test]
1497 async fn update_run_lease_set_with_running_status_attaches_lease() {
1498 let store = InMemoryStore::new();
1499 let run = store
1500 .create_run(new_run_req("test"))
1501 .await
1502 .unwrap()
1503 .into_run();
1504 let expires_at = Utc::now() + TimeDelta::seconds(60);
1505
1506 store
1507 .update_run(
1508 run.id,
1509 RunUpdate {
1510 status: Some(RunStatus::Running),
1511 lease: Some(LeaseUpdate::Set {
1512 worker_id: "worker-1".to_string(),
1513 expires_at,
1514 }),
1515 ..RunUpdate::default()
1516 },
1517 )
1518 .await
1519 .unwrap();
1520
1521 let after = store.get_run(run.id).await.unwrap().unwrap();
1522 assert_eq!(after.status.state, RunStatus::Running);
1523 assert_eq!(after.worker_id.as_deref(), Some("worker-1"));
1524 assert_eq!(after.lease_expires_at, Some(expires_at));
1525 store
1527 .renew_lease(run.id, lease("worker-1", 90).unwrap())
1528 .await
1529 .unwrap();
1530 }
1531
1532 #[tokio::test]
1533 async fn update_run_lease_release_keeps_run_running_without_lease() {
1534 let store = InMemoryStore::new();
1535 store.create_run(new_run_req("test")).await.unwrap();
1536 let picked = store
1537 .pick_next_pending(lease("worker-1", 90))
1538 .await
1539 .unwrap()
1540 .unwrap();
1541
1542 store
1543 .update_run(
1544 picked.id,
1545 RunUpdate {
1546 lease: Some(LeaseUpdate::Release),
1547 ..RunUpdate::default()
1548 },
1549 )
1550 .await
1551 .unwrap();
1552
1553 let after = store.get_run(picked.id).await.unwrap().unwrap();
1554 assert_eq!(after.status.state, RunStatus::Running);
1555 assert!(after.worker_id.is_none());
1556 assert!(after.lease_expires_at.is_none());
1557 let err = store
1558 .renew_lease(picked.id, lease("worker-1", 90).unwrap())
1559 .await
1560 .unwrap_err();
1561 assert!(matches!(err, StoreError::LeaseLost { .. }));
1562 }
1563
1564 #[tokio::test]
1565 async fn update_run_lease_none_leaves_lease_untouched() {
1566 let store = InMemoryStore::new();
1567 store.create_run(new_run_req("test")).await.unwrap();
1568 let picked = store
1569 .pick_next_pending(lease("worker-1", 90))
1570 .await
1571 .unwrap()
1572 .unwrap();
1573
1574 store
1575 .update_run(
1576 picked.id,
1577 RunUpdate {
1578 cost_usd: Some(Decimal::new(150, 2)),
1579 ..RunUpdate::default()
1580 },
1581 )
1582 .await
1583 .unwrap();
1584
1585 let after = store.get_run(picked.id).await.unwrap().unwrap();
1586 assert_eq!(after.worker_id.as_deref(), Some("worker-1"));
1587 assert_eq!(after.lease_expires_at, picked.lease_expires_at);
1588 }
1589
1590 #[tokio::test]
1593 async fn reap_expired_leases_empty_store() {
1594 let store = InMemoryStore::new();
1595 assert!(store.reap_expired_leases(100).await.unwrap().is_empty());
1596 }
1597
1598 #[tokio::test]
1599 async fn reap_expired_leases_requeues_run() {
1600 let store = InMemoryStore::new();
1601 let picked = pick_with_expired_lease(&store, 3).await;
1602
1603 let reaped = store.reap_expired_leases(100).await.unwrap();
1604
1605 assert_eq!(reaped.len(), 1);
1606 assert_eq!(reaped[0].from, RunStatus::Running);
1607 assert_eq!(reaped[0].to, RunStatus::Pending);
1608
1609 let after = store.get_run(picked.id).await.unwrap().unwrap();
1610 assert_eq!(after.status.state, RunStatus::Pending);
1611 assert_eq!(after.retry_count, 0);
1612 assert_eq!(after.lease_recoveries, 1);
1613 assert!(after.error.is_none());
1614 assert!(after.worker_id.is_none());
1615 assert!(after.lease_expires_at.is_none());
1616 }
1617
1618 #[tokio::test]
1619 async fn reap_expired_leases_requeued_run_is_pickable_again() {
1620 let store = InMemoryStore::new();
1621 let picked = pick_with_expired_lease(&store, 3).await;
1622 store.reap_expired_leases(100).await.unwrap();
1623
1624 let repicked = store
1625 .pick_next_pending(lease("worker-2", 90))
1626 .await
1627 .unwrap()
1628 .unwrap();
1629
1630 assert_eq!(repicked.id, picked.id);
1631 assert_eq!(repicked.worker_id.as_deref(), Some("worker-2"));
1632 }
1633
1634 #[tokio::test]
1635 async fn reap_expired_leases_ignores_valid_lease() {
1636 let store = InMemoryStore::new();
1637 store.create_run(new_run_req("test")).await.unwrap();
1638 let picked = store
1639 .pick_next_pending(lease("worker-1", 90))
1640 .await
1641 .unwrap()
1642 .unwrap();
1643
1644 assert!(store.reap_expired_leases(100).await.unwrap().is_empty());
1645
1646 let after = store.get_run(picked.id).await.unwrap().unwrap();
1647 assert_eq!(after.status.state, RunStatus::Running);
1648 assert_eq!(after.retry_count, 0);
1649 }
1650
1651 #[tokio::test]
1652 async fn reap_expired_leases_ignores_run_without_lease() {
1653 let store = InMemoryStore::new();
1654 store.create_run(new_run_req("test")).await.unwrap();
1655 let picked = store.pick_next_pending(None).await.unwrap().unwrap();
1656
1657 assert!(store.reap_expired_leases(100).await.unwrap().is_empty());
1658
1659 let after = store.get_run(picked.id).await.unwrap().unwrap();
1660 assert_eq!(after.status.state, RunStatus::Running);
1661 }
1662
1663 #[tokio::test]
1664 async fn reap_expired_leases_fails_run_when_retries_exhausted() {
1665 let store = InMemoryStore::new();
1666 let picked = pick_with_expired_lease(&store, 0).await;
1667
1668 let reaped = store.reap_expired_leases(100).await.unwrap();
1669
1670 assert_eq!(reaped[0].to, RunStatus::Failed);
1671 let after = store.get_run(picked.id).await.unwrap().unwrap();
1672 assert_eq!(after.status.state, RunStatus::Failed);
1673 assert_eq!(after.error.as_deref(), Some(LEASE_EXPIRED_ERROR));
1674 assert!(after.completed_at.is_some());
1675 }
1676
1677 #[tokio::test]
1678 async fn reap_expired_leases_fails_after_max_retries_recoveries() {
1679 let store = InMemoryStore::new();
1680 let picked = pick_with_expired_lease(&store, 2).await;
1681
1682 for expected in [RunStatus::Pending, RunStatus::Pending, RunStatus::Failed] {
1684 let reaped = store.reap_expired_leases(100).await.unwrap();
1685 assert_eq!(reaped[0].to, expected);
1686 if expected == RunStatus::Pending {
1687 expire_now(&store).await;
1688 }
1689 }
1690
1691 let after = store.get_run(picked.id).await.unwrap().unwrap();
1692 assert_eq!(after.lease_recoveries, 3);
1693 assert_eq!(after.retry_count, 0);
1694 assert_eq!(after.error.as_deref(), Some(LEASE_EXPIRED_ERROR));
1695 }
1696
1697 #[tokio::test]
1698 async fn reap_expired_leases_keeps_the_attempt_number() {
1699 let store = InMemoryStore::new();
1700 let picked = pick_with_expired_lease(&store, 3).await;
1701 let before = store
1702 .create_step(new_step_req(picked.id, "build", 0))
1703 .await
1704 .unwrap();
1705
1706 store.reap_expired_leases(100).await.unwrap();
1707 let repicked = store
1708 .pick_next_pending(lease("worker-2", 90))
1709 .await
1710 .unwrap()
1711 .unwrap();
1712 let after = store
1713 .create_step(new_step_req(repicked.id, "build", 0))
1714 .await
1715 .unwrap();
1716
1717 assert_eq!(before.attempt, 1);
1718 assert_eq!(after.attempt, before.attempt);
1719 assert_eq!(repicked.retry_count, 0);
1720 assert_eq!(repicked.lease_recoveries, 1);
1721 }
1722
1723 #[tokio::test]
1724 async fn reap_expired_leases_counts_apart_from_handler_retries() {
1725 let store = InMemoryStore::new();
1726 let picked = pick_with_expired_lease(&store, 1).await;
1727 store
1728 .update_run(
1729 picked.id,
1730 RunUpdate {
1731 increment_retry: true,
1732 ..RunUpdate::default()
1733 },
1734 )
1735 .await
1736 .unwrap();
1737
1738 let reaped = store.reap_expired_leases(100).await.unwrap();
1741
1742 assert_eq!(reaped[0].to, RunStatus::Pending);
1743 assert_eq!(reaped[0].run.retry_count, 1);
1744 assert_eq!(reaped[0].run.lease_recoveries, 1);
1745 }
1746
1747 #[tokio::test]
1748 async fn reap_expired_leases_respects_limit() {
1749 let store = InMemoryStore::new();
1750 for _ in 0..3 {
1751 pick_with_expired_lease(&store, 3).await;
1752 }
1753
1754 let reaped = store.reap_expired_leases(2).await.unwrap();
1755 assert_eq!(reaped.len(), 2);
1756
1757 let rest = store.reap_expired_leases(100).await.unwrap();
1758 assert_eq!(rest.len(), 1);
1759 }
1760
1761 #[tokio::test]
1764 async fn pick_next_pending_empty_store() {
1765 let store = InMemoryStore::new();
1766 let result = store.pick_next_pending(None).await.unwrap();
1767 assert!(result.is_none());
1768 }
1769
1770 #[tokio::test]
1771 async fn pick_next_pending_returns_oldest_and_transitions_to_running() {
1772 let store = InMemoryStore::new();
1773 let r1 = store
1774 .create_run(new_run_req("first"))
1775 .await
1776 .unwrap()
1777 .into_run();
1778 let _r2 = store
1779 .create_run(new_run_req("second"))
1780 .await
1781 .unwrap()
1782 .into_run();
1783
1784 let picked = store.pick_next_pending(None).await.unwrap().unwrap();
1785 assert_eq!(picked.id, r1.id);
1786 assert_eq!(picked.status.state, RunStatus::Running);
1787 assert!(picked.started_at.is_some());
1788
1789 let fetched = store.get_run(r1.id).await.unwrap().unwrap();
1791 assert_eq!(fetched.status.state, RunStatus::Running);
1792 }
1793
1794 #[tokio::test]
1795 async fn pick_next_pending_skips_non_pending() {
1796 let store = InMemoryStore::new();
1797 let r1 = store.create_run(new_run_req("a")).await.unwrap().into_run();
1798 let r2 = store.create_run(new_run_req("b")).await.unwrap().into_run();
1799
1800 store
1802 .update_run_status(r1.id, RunStatus::Running)
1803 .await
1804 .unwrap();
1805
1806 let picked = store.pick_next_pending(None).await.unwrap().unwrap();
1807 assert_eq!(picked.id, r2.id);
1808 }
1809
1810 #[tokio::test]
1813 async fn create_step_returns_pending() {
1814 let store = InMemoryStore::new();
1815 let run = store
1816 .create_run(new_run_req("test"))
1817 .await
1818 .unwrap()
1819 .into_run();
1820
1821 let step = store
1822 .create_step(NewStep {
1823 run_id: run.id,
1824 trace_id: step_trace_id(run.id, "build", 0),
1825 name: "build".to_string(),
1826 kind: crate::entities::StepKind::Shell,
1827 position: 0,
1828 input: Some(json!({"command": "cargo build"})),
1829 is_error_handler: false,
1830 })
1831 .await
1832 .unwrap();
1833
1834 assert_eq!(step.status.state, StepStatus::Pending);
1835 assert_eq!(step.name, "build");
1836 assert_eq!(step.run_id, run.id);
1837 assert_eq!(step.position, 0);
1838 }
1839
1840 #[tokio::test]
1841 async fn create_step_for_missing_run_returns_error() {
1842 let store = InMemoryStore::new();
1843 let result = store
1844 .create_step(NewStep {
1845 run_id: Uuid::nil(),
1846 trace_id: step_trace_id(Uuid::nil(), "build", 0),
1847 name: "build".to_string(),
1848 kind: crate::entities::StepKind::Shell,
1849 position: 0,
1850 input: None,
1851 is_error_handler: false,
1852 })
1853 .await;
1854 assert!(matches!(result.unwrap_err(), StoreError::RunNotFound(_)));
1855 }
1856
1857 #[tokio::test]
1860 async fn update_step_applies_partial_update() {
1861 let store = InMemoryStore::new();
1862 let run = store
1863 .create_run(new_run_req("test"))
1864 .await
1865 .unwrap()
1866 .into_run();
1867
1868 let step = store
1869 .create_step(NewStep {
1870 run_id: run.id,
1871 trace_id: step_trace_id(run.id, "build", 0),
1872 name: "build".to_string(),
1873 kind: crate::entities::StepKind::Shell,
1874 position: 0,
1875 input: None,
1876 is_error_handler: false,
1877 })
1878 .await
1879 .unwrap();
1880
1881 store
1883 .update_step(
1884 step.id,
1885 StepUpdate {
1886 status: Some(StepStatus::Running),
1887 ..StepUpdate::default()
1888 },
1889 )
1890 .await
1891 .unwrap();
1892
1893 store
1895 .update_step(
1896 step.id,
1897 StepUpdate {
1898 status: Some(StepStatus::Completed),
1899 output: Some(json!({"stdout": "ok"})),
1900 duration_ms: Some(150),
1901 ..StepUpdate::default()
1902 },
1903 )
1904 .await
1905 .unwrap();
1906
1907 let steps = store.list_steps(run.id).await.unwrap();
1908 assert_eq!(steps.len(), 1);
1909 assert_eq!(steps[0].status.state, StepStatus::Completed);
1910 assert_eq!(steps[0].duration_ms, 150);
1911 assert!(steps[0].output.is_some());
1912 }
1913
1914 #[tokio::test]
1915 async fn update_step_not_found() {
1916 let store = InMemoryStore::new();
1917 let result = store.update_step(Uuid::nil(), StepUpdate::default()).await;
1918 assert!(matches!(result.unwrap_err(), StoreError::StepNotFound(_)));
1919 }
1920
1921 #[tokio::test]
1924 async fn list_steps_ordered_by_position() {
1925 let store = InMemoryStore::new();
1926 let run = store
1927 .create_run(new_run_req("test"))
1928 .await
1929 .unwrap()
1930 .into_run();
1931
1932 store
1934 .create_step(NewStep {
1935 run_id: run.id,
1936 trace_id: step_trace_id(run.id, "deploy", 2),
1937 name: "deploy".to_string(),
1938 kind: crate::entities::StepKind::Shell,
1939 position: 2,
1940 input: None,
1941 is_error_handler: false,
1942 })
1943 .await
1944 .unwrap();
1945 store
1946 .create_step(NewStep {
1947 run_id: run.id,
1948 trace_id: step_trace_id(run.id, "build", 0),
1949 name: "build".to_string(),
1950 kind: crate::entities::StepKind::Shell,
1951 position: 0,
1952 input: None,
1953 is_error_handler: false,
1954 })
1955 .await
1956 .unwrap();
1957 store
1958 .create_step(NewStep {
1959 run_id: run.id,
1960 trace_id: step_trace_id(run.id, "test", 1),
1961 name: "test".to_string(),
1962 kind: crate::entities::StepKind::Shell,
1963 position: 1,
1964 input: None,
1965 is_error_handler: false,
1966 })
1967 .await
1968 .unwrap();
1969
1970 let steps = store.list_steps(run.id).await.unwrap();
1971 assert_eq!(steps.len(), 3);
1972 assert_eq!(steps[0].name, "build");
1973 assert_eq!(steps[1].name, "test");
1974 assert_eq!(steps[2].name, "deploy");
1975 }
1976
1977 #[tokio::test]
1978 async fn list_steps_empty_for_run_without_steps() {
1979 let store = InMemoryStore::new();
1980 let run = store
1981 .create_run(new_run_req("test"))
1982 .await
1983 .unwrap()
1984 .into_run();
1985 let steps = store.list_steps(run.id).await.unwrap();
1986 assert!(steps.is_empty());
1987 }
1988
1989 #[tokio::test]
1992 async fn update_run_applies_cost_and_duration() {
1993 let store = InMemoryStore::new();
1994 let run = store
1995 .create_run(new_run_req("test"))
1996 .await
1997 .unwrap()
1998 .into_run();
1999
2000 store
2001 .update_run(
2002 run.id,
2003 RunUpdate {
2004 cost_usd: Some(Decimal::new(123, 2)),
2005 duration_ms: Some(5000),
2006 ..RunUpdate::default()
2007 },
2008 )
2009 .await
2010 .unwrap();
2011
2012 let fetched = store.get_run(run.id).await.unwrap().unwrap();
2013 assert_eq!(fetched.cost_usd, Decimal::new(123, 2));
2014 assert_eq!(fetched.duration_ms, 5000);
2015 }
2016
2017 #[tokio::test]
2018 async fn update_run_increment_retry() {
2019 let store = InMemoryStore::new();
2020 let run = store
2021 .create_run(new_run_req("test"))
2022 .await
2023 .unwrap()
2024 .into_run();
2025 assert_eq!(run.retry_count, 0);
2026
2027 store
2028 .update_run(
2029 run.id,
2030 RunUpdate {
2031 increment_retry: true,
2032 ..RunUpdate::default()
2033 },
2034 )
2035 .await
2036 .unwrap();
2037
2038 let fetched = store.get_run(run.id).await.unwrap().unwrap();
2039 assert_eq!(fetched.retry_count, 1);
2040 }
2041
2042 #[tokio::test]
2043 async fn update_run_not_found() {
2044 let store = InMemoryStore::new();
2045 let result = store.update_run(Uuid::nil(), RunUpdate::default()).await;
2046 assert!(matches!(result.unwrap_err(), StoreError::RunNotFound(_)));
2047 }
2048
2049 #[tokio::test]
2052 async fn concurrent_pick_next_pending_no_double_pick() {
2053 let store = InMemoryStore::new();
2054
2055 for i in 0..10 {
2057 store
2058 .create_run(new_run_req(&format!("wf-{i}")))
2059 .await
2060 .unwrap()
2061 .into_run();
2062 }
2063
2064 let mut handles = Vec::new();
2066 for _ in 0..10 {
2067 let s = store.clone();
2068 handles.push(spawn(async move { s.pick_next_pending(None).await }));
2069 }
2070
2071 let mut picked_ids = Vec::new();
2072 for h in handles {
2073 if let Ok(Ok(Some(run))) = h.await {
2074 picked_ids.push(run.id);
2075 }
2076 }
2077
2078 let unique: std::collections::HashSet<_> = picked_ids.iter().collect();
2080 assert_eq!(unique.len(), picked_ids.len());
2081 }
2082
2083 #[tokio::test]
2086 async fn get_stats_empty_store() {
2087 let store = InMemoryStore::new();
2088 let stats = store.get_stats(RunFilter::default()).await.unwrap();
2089 assert_eq!(stats.total_runs, 0);
2090 assert_eq!(stats.completed_runs, 0);
2091 assert_eq!(stats.failed_runs, 0);
2092 assert_eq!(stats.cancelled_runs, 0);
2093 assert_eq!(stats.active_runs, 0);
2094 assert_eq!(stats.total_cost_usd, Decimal::ZERO);
2095 assert_eq!(stats.total_duration_ms, 0);
2096 }
2097
2098 #[tokio::test]
2099 async fn get_stats_aggregates_counts_and_totals() {
2100 let store = InMemoryStore::new();
2101
2102 let r1 = store
2104 .create_run(new_run_req("wf1"))
2105 .await
2106 .unwrap()
2107 .into_run();
2108 let r2 = store
2109 .create_run(new_run_req("wf2"))
2110 .await
2111 .unwrap()
2112 .into_run();
2113 let r3 = store
2114 .create_run(new_run_req("wf3"))
2115 .await
2116 .unwrap()
2117 .into_run();
2118 let _r4 = store
2119 .create_run(new_run_req("wf4"))
2120 .await
2121 .unwrap()
2122 .into_run();
2123
2124 store
2126 .update_run_status(r1.id, RunStatus::Running)
2127 .await
2128 .unwrap();
2129 store
2130 .update_run_status(r1.id, RunStatus::Completed)
2131 .await
2132 .unwrap();
2133
2134 store
2136 .update_run_status(r2.id, RunStatus::Running)
2137 .await
2138 .unwrap();
2139 store
2140 .update_run_status(r2.id, RunStatus::Failed)
2141 .await
2142 .unwrap();
2143
2144 store
2146 .update_run_status(r3.id, RunStatus::Cancelled)
2147 .await
2148 .unwrap();
2149
2150 store
2154 .update_run(
2155 r1.id,
2156 RunUpdate {
2157 cost_usd: Some(Decimal::new(1000, 2)),
2158 duration_ms: Some(1000),
2159 ..RunUpdate::default()
2160 },
2161 )
2162 .await
2163 .unwrap();
2164
2165 store
2166 .update_run(
2167 r2.id,
2168 RunUpdate {
2169 cost_usd: Some(Decimal::new(500, 2)),
2170 duration_ms: Some(500),
2171 ..RunUpdate::default()
2172 },
2173 )
2174 .await
2175 .unwrap();
2176
2177 let stats = store.get_stats(RunFilter::default()).await.unwrap();
2178 assert_eq!(stats.total_runs, 4);
2179 assert_eq!(stats.completed_runs, 1);
2180 assert_eq!(stats.failed_runs, 1);
2181 assert_eq!(stats.cancelled_runs, 1);
2182 assert_eq!(stats.active_runs, 1); assert_eq!(stats.total_cost_usd, Decimal::new(1500, 2));
2184 assert_eq!(stats.total_duration_ms, 1500);
2185 }
2186
2187 #[tokio::test]
2188 async fn update_run_status_running_to_retrying() {
2189 let store = InMemoryStore::new();
2190 let run = store
2191 .create_run(new_run_req("test"))
2192 .await
2193 .unwrap()
2194 .into_run();
2195
2196 store
2197 .update_run_status(run.id, RunStatus::Running)
2198 .await
2199 .unwrap();
2200
2201 store
2202 .update_run_status(run.id, RunStatus::Retrying)
2203 .await
2204 .unwrap();
2205
2206 let fetched = store.get_run(run.id).await.unwrap().unwrap();
2207 assert_eq!(fetched.status.state, RunStatus::Retrying);
2208 assert!(!fetched.status.state.is_terminal());
2209 assert!(fetched.completed_at.is_none()); }
2211
2212 #[tokio::test]
2213 async fn update_run_status_retrying_to_running_allowed() {
2214 let store = InMemoryStore::new();
2215 let run = store
2216 .create_run(new_run_req("test"))
2217 .await
2218 .unwrap()
2219 .into_run();
2220
2221 store
2222 .update_run_status(run.id, RunStatus::Running)
2223 .await
2224 .unwrap();
2225 store
2226 .update_run_status(run.id, RunStatus::Retrying)
2227 .await
2228 .unwrap();
2229
2230 store
2232 .update_run_status(run.id, RunStatus::Running)
2233 .await
2234 .unwrap();
2235
2236 let fetched = store.get_run(run.id).await.unwrap().unwrap();
2237 assert_eq!(fetched.status.state, RunStatus::Running);
2238 }
2239
2240 #[tokio::test]
2241 async fn update_run_with_invalid_status_transition_errors() {
2242 let store = InMemoryStore::new();
2243 let run = store
2244 .create_run(new_run_req("test"))
2245 .await
2246 .unwrap()
2247 .into_run();
2248
2249 let result = store
2251 .update_run(
2252 run.id,
2253 RunUpdate {
2254 status: Some(RunStatus::Completed), ..RunUpdate::default()
2256 },
2257 )
2258 .await;
2259
2260 assert!(result.is_err());
2261 }
2262
2263 #[tokio::test]
2264 async fn create_step_with_complex_input() {
2265 let store = InMemoryStore::new();
2266 let run = store
2267 .create_run(new_run_req("test"))
2268 .await
2269 .unwrap()
2270 .into_run();
2271
2272 let complex_input = json!({
2273 "command": "cargo build",
2274 "env": {
2275 "RUST_LOG": "debug",
2276 "CUSTOM": "value"
2277 },
2278 "timeout": 60,
2279 "retry_policy": {
2280 "max_attempts": 3,
2281 "backoff": "exponential"
2282 }
2283 });
2284
2285 let step = store
2286 .create_step(NewStep {
2287 run_id: run.id,
2288 trace_id: step_trace_id(run.id, "build", 0),
2289 name: "build".to_string(),
2290 kind: crate::entities::StepKind::Agent,
2291 position: 0,
2292 input: Some(complex_input.clone()),
2293 is_error_handler: false,
2294 })
2295 .await
2296 .unwrap();
2297
2298 assert_eq!(step.input, Some(complex_input));
2299 }
2300
2301 #[tokio::test]
2302 async fn update_step_with_error_message() {
2303 let store = InMemoryStore::new();
2304 let run = store
2305 .create_run(new_run_req("test"))
2306 .await
2307 .unwrap()
2308 .into_run();
2309
2310 let step = store
2311 .create_step(NewStep {
2312 run_id: run.id,
2313 trace_id: step_trace_id(run.id, "build", 0),
2314 name: "build".to_string(),
2315 kind: crate::entities::StepKind::Shell,
2316 position: 0,
2317 input: None,
2318 is_error_handler: false,
2319 })
2320 .await
2321 .unwrap();
2322
2323 store
2324 .update_step(
2325 step.id,
2326 StepUpdate {
2327 status: Some(StepStatus::Running),
2328 ..StepUpdate::default()
2329 },
2330 )
2331 .await
2332 .unwrap();
2333
2334 store
2335 .update_step(
2336 step.id,
2337 StepUpdate {
2338 status: Some(StepStatus::Failed),
2339 error: Some("Connection timeout after 30s".to_string()),
2340 duration_ms: Some(30000),
2341 ..StepUpdate::default()
2342 },
2343 )
2344 .await
2345 .unwrap();
2346
2347 let steps = store.list_steps(run.id).await.unwrap();
2348 assert_eq!(steps[0].status.state, StepStatus::Failed);
2349 assert_eq!(
2350 steps[0].error,
2351 Some("Connection timeout after 30s".to_string())
2352 );
2353 assert_eq!(steps[0].duration_ms, 30000);
2354 }
2355
2356 #[tokio::test]
2357 async fn list_steps_for_nonexistent_run_returns_empty() {
2358 let store = InMemoryStore::new();
2359 let steps = store.list_steps(Uuid::nil()).await.unwrap();
2360 assert!(steps.is_empty());
2361 }
2362
2363 #[tokio::test]
2364 async fn update_step_pending_to_skipped() {
2365 let store = InMemoryStore::new();
2366 let run = store
2367 .create_run(new_run_req("test"))
2368 .await
2369 .unwrap()
2370 .into_run();
2371
2372 let step = store
2373 .create_step(NewStep {
2374 run_id: run.id,
2375 trace_id: step_trace_id(run.id, "build", 0),
2376 name: "build".to_string(),
2377 kind: crate::entities::StepKind::Shell,
2378 position: 0,
2379 input: None,
2380 is_error_handler: false,
2381 })
2382 .await
2383 .unwrap();
2384
2385 store
2387 .update_step(
2388 step.id,
2389 StepUpdate {
2390 status: Some(StepStatus::Skipped),
2391 ..StepUpdate::default()
2392 },
2393 )
2394 .await
2395 .unwrap();
2396
2397 let steps = store.list_steps(run.id).await.unwrap();
2398 assert_eq!(steps[0].status.state, StepStatus::Skipped);
2399 }
2400
2401 #[tokio::test]
2402 async fn list_runs_with_combined_filters() {
2403 let store = InMemoryStore::new();
2404
2405 let r1 = store
2406 .create_run(new_run_req("deploy"))
2407 .await
2408 .unwrap()
2409 .into_run();
2410 let r2 = store
2411 .create_run(new_run_req("deploy"))
2412 .await
2413 .unwrap()
2414 .into_run();
2415 let _r3 = store
2416 .create_run(new_run_req("test"))
2417 .await
2418 .unwrap()
2419 .into_run();
2420
2421 store
2423 .update_run_status(r1.id, RunStatus::Running)
2424 .await
2425 .unwrap();
2426 store
2427 .update_run_status(r1.id, RunStatus::Completed)
2428 .await
2429 .unwrap();
2430
2431 store
2433 .update_run_status(r2.id, RunStatus::Running)
2434 .await
2435 .unwrap();
2436
2437 let filter = RunFilter {
2439 workflow_name: Some("deploy".to_string()),
2440 status: Some(RunStatus::Running),
2441 ..RunFilter::default()
2442 };
2443
2444 let page = store.list_runs(filter, 1, 100).await.unwrap();
2445 assert_eq!(page.total, 1);
2446 assert_eq!(page.items[0].id, r2.id);
2447 }
2448
2449 #[tokio::test]
2450 async fn list_runs_workflow_filter_is_case_insensitive_partial_match() {
2451 let store = InMemoryStore::new();
2452 store
2453 .create_run(new_run_req("weather-report"))
2454 .await
2455 .unwrap()
2456 .into_run();
2457 store
2458 .create_run(new_run_req("deploy-prod"))
2459 .await
2460 .unwrap()
2461 .into_run();
2462
2463 let filter = RunFilter {
2465 workflow_name: Some("weather".to_string()),
2466 ..RunFilter::default()
2467 };
2468 let page = store.list_runs(filter, 1, 100).await.unwrap();
2469 assert_eq!(page.total, 1);
2470 assert_eq!(page.items[0].workflow_name, "weather-report");
2471
2472 let filter = RunFilter {
2474 workflow_name: Some("Weather-REPORT".to_string()),
2475 ..RunFilter::default()
2476 };
2477 let page = store.list_runs(filter, 1, 100).await.unwrap();
2478 assert_eq!(page.total, 1);
2479 assert_eq!(page.items[0].workflow_name, "weather-report");
2480
2481 let filter = RunFilter {
2483 workflow_name: Some("report".to_string()),
2484 ..RunFilter::default()
2485 };
2486 let page = store.list_runs(filter, 1, 100).await.unwrap();
2487 assert_eq!(page.total, 1);
2488 assert_eq!(page.items[0].workflow_name, "weather-report");
2489
2490 let filter = RunFilter {
2492 workflow_name: Some("build".to_string()),
2493 ..RunFilter::default()
2494 };
2495 let page = store.list_runs(filter, 1, 100).await.unwrap();
2496 assert_eq!(page.total, 0);
2497 }
2498
2499 #[tokio::test]
2500 async fn list_runs_has_steps_true_only_filters_completed_and_cancelled() {
2501 let store = InMemoryStore::new();
2502 let run_with = create_terminal_run(&store, "with-steps", RunStatus::Completed).await;
2503 let _run_without = create_terminal_run(&store, "without-steps", RunStatus::Completed).await;
2504
2505 store
2506 .create_step(NewStep {
2507 run_id: run_with.id,
2508 trace_id: step_trace_id(run_with.id, "build", 0),
2509 name: "build".to_string(),
2510 kind: crate::entities::StepKind::Shell,
2511 position: 0,
2512 input: None,
2513 is_error_handler: false,
2514 })
2515 .await
2516 .unwrap();
2517
2518 let filter = RunFilter {
2519 has_steps: Some(true),
2520 ..RunFilter::default()
2521 };
2522 let page = store.list_runs(filter, 1, 100).await.unwrap();
2523 assert_eq!(page.total, 1);
2524 assert_eq!(page.items[0].id, run_with.id);
2525 }
2526
2527 #[tokio::test]
2528 async fn list_runs_has_steps_false_only_filters_completed_and_cancelled() {
2529 let store = InMemoryStore::new();
2530 let run_with = create_terminal_run(&store, "with-steps", RunStatus::Cancelled).await;
2531 let run_without = create_terminal_run(&store, "without-steps", RunStatus::Cancelled).await;
2532
2533 store
2534 .create_step(NewStep {
2535 run_id: run_with.id,
2536 trace_id: step_trace_id(run_with.id, "build", 0),
2537 name: "build".to_string(),
2538 kind: crate::entities::StepKind::Shell,
2539 position: 0,
2540 input: None,
2541 is_error_handler: false,
2542 })
2543 .await
2544 .unwrap();
2545
2546 let filter = RunFilter {
2547 has_steps: Some(false),
2548 ..RunFilter::default()
2549 };
2550 let page = store.list_runs(filter, 1, 100).await.unwrap();
2551 assert_eq!(page.total, 1);
2552 assert_eq!(page.items[0].id, run_without.id);
2553 }
2554
2555 #[tokio::test]
2556 async fn list_runs_has_steps_none_returns_all() {
2557 let store = InMemoryStore::new();
2558 let run_with = store
2559 .create_run(new_run_req("with-steps"))
2560 .await
2561 .unwrap()
2562 .into_run();
2563 let _run_without = store
2564 .create_run(new_run_req("without-steps"))
2565 .await
2566 .unwrap()
2567 .into_run();
2568
2569 store
2570 .create_step(NewStep {
2571 run_id: run_with.id,
2572 trace_id: step_trace_id(run_with.id, "build", 0),
2573 name: "build".to_string(),
2574 kind: crate::entities::StepKind::Shell,
2575 position: 0,
2576 input: None,
2577 is_error_handler: false,
2578 })
2579 .await
2580 .unwrap();
2581
2582 let filter = RunFilter {
2583 has_steps: None,
2584 ..RunFilter::default()
2585 };
2586 let page = store.list_runs(filter, 1, 100).await.unwrap();
2587 assert_eq!(page.total, 2);
2588 }
2589
2590 #[tokio::test]
2591 async fn list_runs_has_steps_true_does_not_filter_non_terminal_runs() {
2592 let store = InMemoryStore::new();
2593 let pending_run = store
2594 .create_run(new_run_req("pending-empty"))
2595 .await
2596 .unwrap()
2597 .into_run();
2598 let running_run = store
2599 .create_run(new_run_req("running-empty"))
2600 .await
2601 .unwrap()
2602 .into_run();
2603 store
2604 .update_run_status(running_run.id, RunStatus::Running)
2605 .await
2606 .unwrap();
2607
2608 let filter = RunFilter {
2609 has_steps: Some(true),
2610 ..RunFilter::default()
2611 };
2612 let page = store.list_runs(filter, 1, 100).await.unwrap();
2613 assert_eq!(page.total, 2);
2614 let ids: Vec<_> = page.items.iter().map(|r| r.id).collect();
2615 assert!(ids.contains(&pending_run.id));
2616 assert!(ids.contains(&running_run.id));
2617 }
2618
2619 #[tokio::test]
2620 async fn get_stats_with_mixed_active_statuses() {
2621 let store = InMemoryStore::new();
2622
2623 let _r1 = store
2624 .create_run(new_run_req("wf"))
2625 .await
2626 .unwrap()
2627 .into_run(); let r2 = store
2629 .create_run(new_run_req("wf"))
2630 .await
2631 .unwrap()
2632 .into_run();
2633 let r3 = store
2634 .create_run(new_run_req("wf"))
2635 .await
2636 .unwrap()
2637 .into_run();
2638
2639 store
2640 .update_run_status(r2.id, RunStatus::Running)
2641 .await
2642 .unwrap();
2643 store
2644 .update_run_status(r3.id, RunStatus::Running)
2645 .await
2646 .unwrap();
2647 store
2648 .update_run_status(r3.id, RunStatus::Retrying)
2649 .await
2650 .unwrap();
2651
2652 let r4 = store
2653 .create_run(new_run_req("wf"))
2654 .await
2655 .unwrap()
2656 .into_run();
2657 store
2658 .update_run_status(r4.id, RunStatus::Running)
2659 .await
2660 .unwrap();
2661 store
2662 .update_run_status(r4.id, RunStatus::AwaitingApproval)
2663 .await
2664 .unwrap();
2665
2666 let r5 = store
2667 .create_run(new_run_req("wf"))
2668 .await
2669 .unwrap()
2670 .into_run();
2671 store
2672 .update_run_status(r5.id, RunStatus::Running)
2673 .await
2674 .unwrap();
2675 store
2676 .update_run_status(r5.id, RunStatus::Sleeping)
2677 .await
2678 .unwrap();
2679
2680 let stats = store.get_stats(RunFilter::default()).await.unwrap();
2681 assert_eq!(stats.active_runs, 5);
2683 assert_eq!(stats.awaiting_approval_runs, 1);
2684 }
2685
2686 #[tokio::test]
2687 async fn run_with_different_trigger_kinds() {
2688 let store = InMemoryStore::new();
2689
2690 let r1 = store
2691 .create_run(NewRun {
2692 created_by: None,
2693 workflow_name: "test".to_string(),
2694 trigger: TriggerKind::Manual,
2695 payload: json!({}),
2696 max_retries: 1,
2697 handler_version: None,
2698 labels: HashMap::new(),
2699 scheduled_at: None,
2700 idempotency_key: None,
2701 concurrency_key: None,
2702 concurrency_limits: Vec::new(),
2703 max_cost_usd: None,
2704 })
2705 .await
2706 .unwrap()
2707 .into_run();
2708
2709 let r2 = store
2710 .create_run(NewRun {
2711 created_by: None,
2712 workflow_name: "test".to_string(),
2713 trigger: TriggerKind::Webhook {
2714 path: "/hooks/github".to_string(),
2715 },
2716 payload: json!({}),
2717 max_retries: 1,
2718 handler_version: None,
2719 labels: HashMap::new(),
2720 scheduled_at: None,
2721 idempotency_key: None,
2722 concurrency_key: None,
2723 concurrency_limits: Vec::new(),
2724 max_cost_usd: None,
2725 })
2726 .await
2727 .unwrap()
2728 .into_run();
2729
2730 let r3 = store
2731 .create_run(NewRun {
2732 created_by: None,
2733 workflow_name: "test".to_string(),
2734 trigger: TriggerKind::Cron {
2735 schedule: "0 0 * * *".to_string(),
2736 },
2737 payload: json!({}),
2738 max_retries: 1,
2739 handler_version: None,
2740 labels: HashMap::new(),
2741 scheduled_at: None,
2742 idempotency_key: None,
2743 concurrency_key: None,
2744 concurrency_limits: Vec::new(),
2745 max_cost_usd: None,
2746 })
2747 .await
2748 .unwrap()
2749 .into_run();
2750
2751 let r4 = store
2752 .create_run(NewRun {
2753 created_by: None,
2754 workflow_name: "test".to_string(),
2755 trigger: TriggerKind::Api,
2756 payload: json!({}),
2757 max_retries: 1,
2758 handler_version: None,
2759 labels: HashMap::new(),
2760 scheduled_at: None,
2761 idempotency_key: None,
2762 concurrency_key: None,
2763 concurrency_limits: Vec::new(),
2764 max_cost_usd: None,
2765 })
2766 .await
2767 .unwrap()
2768 .into_run();
2769
2770 let r5 = store
2771 .create_run(NewRun {
2772 created_by: None,
2773 workflow_name: "test".to_string(),
2774 trigger: TriggerKind::Retry {
2775 parent_run_id: Uuid::nil(),
2776 },
2777 payload: json!({}),
2778 max_retries: 1,
2779 handler_version: None,
2780 labels: HashMap::new(),
2781 scheduled_at: None,
2782 idempotency_key: None,
2783 concurrency_key: None,
2784 concurrency_limits: Vec::new(),
2785 max_cost_usd: None,
2786 })
2787 .await
2788 .unwrap()
2789 .into_run();
2790
2791 assert_eq!(r1.trigger, TriggerKind::Manual);
2792 assert!(matches!(r2.trigger, TriggerKind::Webhook { .. }));
2793 assert!(matches!(r3.trigger, TriggerKind::Cron { .. }));
2794 assert_eq!(r4.trigger, TriggerKind::Api);
2795 assert!(matches!(r5.trigger, TriggerKind::Retry { .. }));
2796 }
2797
2798 #[tokio::test]
2801 async fn create_step_dependencies_stores_dependencies() {
2802 let store = InMemoryStore::new();
2803 let run = store
2804 .create_run(new_run_req("test"))
2805 .await
2806 .unwrap()
2807 .into_run();
2808
2809 let step1 = store
2810 .create_step(NewStep {
2811 run_id: run.id,
2812 trace_id: step_trace_id(run.id, "step1", 0),
2813 name: "step1".to_string(),
2814 kind: crate::entities::StepKind::Shell,
2815 position: 0,
2816 input: None,
2817 is_error_handler: false,
2818 })
2819 .await
2820 .unwrap();
2821
2822 let step2 = store
2823 .create_step(NewStep {
2824 run_id: run.id,
2825 trace_id: step_trace_id(run.id, "step2", 1),
2826 name: "step2".to_string(),
2827 kind: crate::entities::StepKind::Shell,
2828 position: 1,
2829 input: None,
2830 is_error_handler: false,
2831 })
2832 .await
2833 .unwrap();
2834
2835 let result = store
2836 .create_step_dependencies(vec![NewStepDependency {
2837 step_id: step2.id,
2838 depends_on: step1.id,
2839 }])
2840 .await;
2841
2842 assert!(result.is_ok());
2843
2844 let deps = store.list_step_dependencies(run.id).await.unwrap();
2845 assert_eq!(deps.len(), 1);
2846 assert_eq!(deps[0].step_id, step2.id);
2847 assert_eq!(deps[0].depends_on, step1.id);
2848 }
2849
2850 #[tokio::test]
2851 async fn create_step_dependencies_duplicate_dependencies_are_idempotent() {
2852 let store = InMemoryStore::new();
2853 let run = store
2854 .create_run(new_run_req("test"))
2855 .await
2856 .unwrap()
2857 .into_run();
2858
2859 let step1 = store
2860 .create_step(NewStep {
2861 run_id: run.id,
2862 trace_id: step_trace_id(run.id, "step1", 0),
2863 name: "step1".to_string(),
2864 kind: crate::entities::StepKind::Shell,
2865 position: 0,
2866 input: None,
2867 is_error_handler: false,
2868 })
2869 .await
2870 .unwrap();
2871
2872 let step2 = store
2873 .create_step(NewStep {
2874 run_id: run.id,
2875 trace_id: step_trace_id(run.id, "step2", 1),
2876 name: "step2".to_string(),
2877 kind: crate::entities::StepKind::Shell,
2878 position: 1,
2879 input: None,
2880 is_error_handler: false,
2881 })
2882 .await
2883 .unwrap();
2884
2885 let dep = NewStepDependency {
2886 step_id: step2.id,
2887 depends_on: step1.id,
2888 };
2889
2890 store
2891 .create_step_dependencies(vec![dep.clone()])
2892 .await
2893 .unwrap();
2894 store.create_step_dependencies(vec![dep]).await.unwrap();
2895
2896 let deps = store.list_step_dependencies(run.id).await.unwrap();
2897 assert_eq!(deps.len(), 1);
2898 }
2899
2900 #[tokio::test]
2901 async fn create_step_dependencies_missing_step_id_returns_error() {
2902 let store = InMemoryStore::new();
2903 let run = store
2904 .create_run(new_run_req("test"))
2905 .await
2906 .unwrap()
2907 .into_run();
2908
2909 let step1 = store
2910 .create_step(NewStep {
2911 run_id: run.id,
2912 trace_id: step_trace_id(run.id, "step1", 0),
2913 name: "step1".to_string(),
2914 kind: crate::entities::StepKind::Shell,
2915 position: 0,
2916 input: None,
2917 is_error_handler: false,
2918 })
2919 .await
2920 .unwrap();
2921
2922 let result = store
2923 .create_step_dependencies(vec![NewStepDependency {
2924 step_id: Uuid::nil(),
2925 depends_on: step1.id,
2926 }])
2927 .await;
2928
2929 assert!(matches!(result.unwrap_err(), StoreError::StepNotFound(_)));
2930 }
2931
2932 #[tokio::test]
2933 async fn create_step_dependencies_missing_depends_on_returns_error() {
2934 let store = InMemoryStore::new();
2935 let run = store
2936 .create_run(new_run_req("test"))
2937 .await
2938 .unwrap()
2939 .into_run();
2940
2941 let step1 = store
2942 .create_step(NewStep {
2943 run_id: run.id,
2944 trace_id: step_trace_id(run.id, "step1", 0),
2945 name: "step1".to_string(),
2946 kind: crate::entities::StepKind::Shell,
2947 position: 0,
2948 input: None,
2949 is_error_handler: false,
2950 })
2951 .await
2952 .unwrap();
2953
2954 let result = store
2955 .create_step_dependencies(vec![NewStepDependency {
2956 step_id: step1.id,
2957 depends_on: Uuid::nil(),
2958 }])
2959 .await;
2960
2961 assert!(matches!(result.unwrap_err(), StoreError::StepNotFound(_)));
2962 }
2963
2964 #[tokio::test]
2965 async fn create_step_dependencies_multiple_dependencies() {
2966 let store = InMemoryStore::new();
2967 let run = store
2968 .create_run(new_run_req("test"))
2969 .await
2970 .unwrap()
2971 .into_run();
2972
2973 let step1 = store
2974 .create_step(NewStep {
2975 run_id: run.id,
2976 trace_id: step_trace_id(run.id, "step1", 0),
2977 name: "step1".to_string(),
2978 kind: crate::entities::StepKind::Shell,
2979 position: 0,
2980 input: None,
2981 is_error_handler: false,
2982 })
2983 .await
2984 .unwrap();
2985
2986 let step2 = store
2987 .create_step(NewStep {
2988 run_id: run.id,
2989 trace_id: step_trace_id(run.id, "step2", 1),
2990 name: "step2".to_string(),
2991 kind: crate::entities::StepKind::Shell,
2992 position: 1,
2993 input: None,
2994 is_error_handler: false,
2995 })
2996 .await
2997 .unwrap();
2998
2999 let step3 = store
3000 .create_step(NewStep {
3001 run_id: run.id,
3002 trace_id: step_trace_id(run.id, "step3", 2),
3003 name: "step3".to_string(),
3004 kind: crate::entities::StepKind::Shell,
3005 position: 2,
3006 input: None,
3007 is_error_handler: false,
3008 })
3009 .await
3010 .unwrap();
3011
3012 let result = store
3013 .create_step_dependencies(vec![
3014 NewStepDependency {
3015 step_id: step2.id,
3016 depends_on: step1.id,
3017 },
3018 NewStepDependency {
3019 step_id: step3.id,
3020 depends_on: step2.id,
3021 },
3022 ])
3023 .await;
3024
3025 assert!(result.is_ok());
3026
3027 let deps = store.list_step_dependencies(run.id).await.unwrap();
3028 assert_eq!(deps.len(), 2);
3029 }
3030
3031 #[tokio::test]
3034 async fn list_step_dependencies_empty_for_run_with_no_dependencies() {
3035 let store = InMemoryStore::new();
3036 let run = store
3037 .create_run(new_run_req("test"))
3038 .await
3039 .unwrap()
3040 .into_run();
3041
3042 store
3043 .create_step(NewStep {
3044 run_id: run.id,
3045 trace_id: step_trace_id(run.id, "step1", 0),
3046 name: "step1".to_string(),
3047 kind: crate::entities::StepKind::Shell,
3048 position: 0,
3049 input: None,
3050 is_error_handler: false,
3051 })
3052 .await
3053 .unwrap();
3054
3055 let deps = store.list_step_dependencies(run.id).await.unwrap();
3056 assert!(deps.is_empty());
3057 }
3058
3059 #[tokio::test]
3060 async fn list_step_dependencies_returns_only_deps_for_given_run() {
3061 let store = InMemoryStore::new();
3062 let run1 = store
3063 .create_run(new_run_req("test1"))
3064 .await
3065 .unwrap()
3066 .into_run();
3067 let run2 = store
3068 .create_run(new_run_req("test2"))
3069 .await
3070 .unwrap()
3071 .into_run();
3072
3073 let step1_run1 = store
3074 .create_step(NewStep {
3075 run_id: run1.id,
3076 trace_id: step_trace_id(run1.id, "step1", 0),
3077 name: "step1".to_string(),
3078 kind: crate::entities::StepKind::Shell,
3079 position: 0,
3080 input: None,
3081 is_error_handler: false,
3082 })
3083 .await
3084 .unwrap();
3085
3086 let step2_run1 = store
3087 .create_step(NewStep {
3088 run_id: run1.id,
3089 trace_id: step_trace_id(run1.id, "step2", 1),
3090 name: "step2".to_string(),
3091 kind: crate::entities::StepKind::Shell,
3092 position: 1,
3093 input: None,
3094 is_error_handler: false,
3095 })
3096 .await
3097 .unwrap();
3098
3099 let step1_run2 = store
3100 .create_step(NewStep {
3101 run_id: run2.id,
3102 trace_id: step_trace_id(run2.id, "step1", 0),
3103 name: "step1".to_string(),
3104 kind: crate::entities::StepKind::Shell,
3105 position: 0,
3106 input: None,
3107 is_error_handler: false,
3108 })
3109 .await
3110 .unwrap();
3111
3112 let step2_run2 = store
3113 .create_step(NewStep {
3114 run_id: run2.id,
3115 trace_id: step_trace_id(run2.id, "step2", 1),
3116 name: "step2".to_string(),
3117 kind: crate::entities::StepKind::Shell,
3118 position: 1,
3119 input: None,
3120 is_error_handler: false,
3121 })
3122 .await
3123 .unwrap();
3124
3125 store
3126 .create_step_dependencies(vec![
3127 NewStepDependency {
3128 step_id: step2_run1.id,
3129 depends_on: step1_run1.id,
3130 },
3131 NewStepDependency {
3132 step_id: step2_run2.id,
3133 depends_on: step1_run2.id,
3134 },
3135 ])
3136 .await
3137 .unwrap();
3138
3139 let deps_run1 = store.list_step_dependencies(run1.id).await.unwrap();
3140 let deps_run2 = store.list_step_dependencies(run2.id).await.unwrap();
3141
3142 assert_eq!(deps_run1.len(), 1);
3143 assert_eq!(deps_run1[0].step_id, step2_run1.id);
3144 assert_eq!(deps_run1[0].depends_on, step1_run1.id);
3145
3146 assert_eq!(deps_run2.len(), 1);
3147 assert_eq!(deps_run2[0].step_id, step2_run2.id);
3148 assert_eq!(deps_run2[0].depends_on, step1_run2.id);
3149 }
3150
3151 #[tokio::test]
3152 async fn list_step_dependencies_returns_empty_for_nonexistent_run() {
3153 let store = InMemoryStore::new();
3154 let deps = store.list_step_dependencies(Uuid::nil()).await.unwrap();
3155 assert!(deps.is_empty());
3156 }
3157
3158 #[tokio::test]
3159 async fn list_step_dependencies_sorted_by_created_at() {
3160 let store = InMemoryStore::new();
3161 let run = store
3162 .create_run(new_run_req("test"))
3163 .await
3164 .unwrap()
3165 .into_run();
3166
3167 let step1 = store
3168 .create_step(NewStep {
3169 run_id: run.id,
3170 trace_id: step_trace_id(run.id, "step1", 0),
3171 name: "step1".to_string(),
3172 kind: crate::entities::StepKind::Shell,
3173 position: 0,
3174 input: None,
3175 is_error_handler: false,
3176 })
3177 .await
3178 .unwrap();
3179
3180 let step2 = store
3181 .create_step(NewStep {
3182 run_id: run.id,
3183 trace_id: step_trace_id(run.id, "step2", 1),
3184 name: "step2".to_string(),
3185 kind: crate::entities::StepKind::Shell,
3186 position: 1,
3187 input: None,
3188 is_error_handler: false,
3189 })
3190 .await
3191 .unwrap();
3192
3193 let step3 = store
3194 .create_step(NewStep {
3195 run_id: run.id,
3196 trace_id: step_trace_id(run.id, "step3", 2),
3197 name: "step3".to_string(),
3198 kind: crate::entities::StepKind::Shell,
3199 position: 2,
3200 input: None,
3201 is_error_handler: false,
3202 })
3203 .await
3204 .unwrap();
3205
3206 store
3207 .create_step_dependencies(vec![NewStepDependency {
3208 step_id: step2.id,
3209 depends_on: step1.id,
3210 }])
3211 .await
3212 .unwrap();
3213
3214 store
3215 .create_step_dependencies(vec![NewStepDependency {
3216 step_id: step3.id,
3217 depends_on: step1.id,
3218 }])
3219 .await
3220 .unwrap();
3221
3222 let deps = store.list_step_dependencies(run.id).await.unwrap();
3223 assert_eq!(deps.len(), 2);
3224 assert!(deps[0].created_at <= deps[1].created_at);
3225 }
3226
3227 #[tokio::test]
3230 async fn update_run_returning_applies_and_returns() {
3231 let store = InMemoryStore::new();
3232 let run = store
3233 .create_run(new_run_req("test"))
3234 .await
3235 .unwrap()
3236 .into_run();
3237
3238 store
3240 .update_run_status(run.id, RunStatus::Running)
3241 .await
3242 .unwrap();
3243
3244 let updated = store
3245 .update_run_returning(
3246 run.id,
3247 RunUpdate {
3248 status: Some(RunStatus::Completed),
3249 cost_usd: Some(Decimal::new(4200, 2)),
3250 duration_ms: Some(1500),
3251 ..RunUpdate::default()
3252 },
3253 )
3254 .await
3255 .unwrap();
3256
3257 assert_eq!(updated.id, run.id);
3258 assert_eq!(updated.status.state, RunStatus::Completed);
3259 assert_eq!(updated.cost_usd, Decimal::new(4200, 2));
3260 assert_eq!(updated.duration_ms, 1500);
3261 assert!(updated.completed_at.is_some());
3262 }
3263
3264 #[tokio::test]
3265 async fn update_run_returning_not_found() {
3266 let store = InMemoryStore::new();
3267 let result = store
3268 .update_run_returning(
3269 Uuid::nil(),
3270 RunUpdate {
3271 status: Some(RunStatus::Running),
3272 ..RunUpdate::default()
3273 },
3274 )
3275 .await;
3276
3277 assert!(matches!(result, Err(StoreError::RunNotFound(_))));
3278 }
3279
3280 #[tokio::test]
3281 async fn update_run_returning_invalid_transition() {
3282 let store = InMemoryStore::new();
3283 let run = store
3284 .create_run(new_run_req("test"))
3285 .await
3286 .unwrap()
3287 .into_run();
3288
3289 let result = store
3290 .update_run_returning(
3291 run.id,
3292 RunUpdate {
3293 status: Some(RunStatus::Completed),
3294 ..RunUpdate::default()
3295 },
3296 )
3297 .await;
3298
3299 assert!(matches!(result, Err(StoreError::InvalidTransition { .. })));
3300 }
3301
3302 #[tokio::test]
3305 async fn create_step_stamps_the_current_attempt() {
3306 let store = InMemoryStore::new();
3307 let run = store
3308 .create_run(new_run_req("retry-wf"))
3309 .await
3310 .unwrap()
3311 .into_run();
3312
3313 let first = store
3314 .create_step(new_step_req(run.id, "build", 0))
3315 .await
3316 .unwrap();
3317 assert_eq!(first.attempt, 1);
3318
3319 store
3320 .update_run_status(run.id, RunStatus::Running)
3321 .await
3322 .unwrap();
3323 store
3324 .update_run(
3325 run.id,
3326 RunUpdate {
3327 status: Some(RunStatus::Retrying),
3328 increment_retry: true,
3329 ..RunUpdate::default()
3330 },
3331 )
3332 .await
3333 .unwrap();
3334
3335 let second = store
3336 .create_step(new_step_req(run.id, "build", 0))
3337 .await
3338 .unwrap();
3339 assert_eq!(second.attempt, 2);
3340 }
3341
3342 #[tokio::test]
3343 async fn pick_next_pending_ignores_retrying_run_before_its_backoff() {
3344 let store = InMemoryStore::new();
3345 let run = store
3346 .create_run(new_run_req("retry-wf"))
3347 .await
3348 .unwrap()
3349 .into_run();
3350
3351 store
3352 .update_run_status(run.id, RunStatus::Running)
3353 .await
3354 .unwrap();
3355 store
3356 .update_run(
3357 run.id,
3358 RunUpdate {
3359 status: Some(RunStatus::Retrying),
3360 increment_retry: true,
3361 scheduled_at: Some(Utc::now() + TimeDelta::seconds(60)),
3362 ..RunUpdate::default()
3363 },
3364 )
3365 .await
3366 .unwrap();
3367
3368 assert!(store.pick_next_pending(None).await.unwrap().is_none());
3369 }
3370
3371 #[tokio::test]
3372 async fn pick_next_pending_resumes_retrying_run_after_its_backoff() {
3373 let store = InMemoryStore::new();
3374 let run = store
3375 .create_run(new_run_req("retry-wf"))
3376 .await
3377 .unwrap()
3378 .into_run();
3379
3380 store
3381 .update_run_status(run.id, RunStatus::Running)
3382 .await
3383 .unwrap();
3384 store
3385 .update_run(
3386 run.id,
3387 RunUpdate {
3388 status: Some(RunStatus::Retrying),
3389 increment_retry: true,
3390 scheduled_at: Some(Utc::now() - TimeDelta::seconds(1)),
3391 ..RunUpdate::default()
3392 },
3393 )
3394 .await
3395 .unwrap();
3396
3397 let picked = store.pick_next_pending(None).await.unwrap().unwrap();
3398 assert_eq!(picked.id, run.id);
3399 assert_eq!(picked.status.state, RunStatus::Running);
3400 assert_eq!(picked.retry_count, 1);
3401 }
3402
3403 #[tokio::test]
3404 async fn update_run_persists_scheduled_at() {
3405 let store = InMemoryStore::new();
3406 let run = store
3407 .create_run(new_run_req("test"))
3408 .await
3409 .unwrap()
3410 .into_run();
3411 let when = Utc::now() + TimeDelta::seconds(30);
3412
3413 store
3414 .update_run(
3415 run.id,
3416 RunUpdate {
3417 scheduled_at: Some(when),
3418 ..RunUpdate::default()
3419 },
3420 )
3421 .await
3422 .unwrap();
3423
3424 let fetched = store.get_run(run.id).await.unwrap().unwrap();
3425 assert_eq!(fetched.scheduled_at, Some(when));
3426 }
3427
3428 #[tokio::test]
3431 async fn new_run_has_no_output() {
3432 let store = InMemoryStore::new();
3433 let run = store
3434 .create_run(new_run_req("test"))
3435 .await
3436 .unwrap()
3437 .into_run();
3438
3439 assert!(run.output.is_none());
3440 let fetched = store.get_run(run.id).await.unwrap().unwrap();
3441 assert!(fetched.output.is_none());
3442 }
3443
3444 #[tokio::test]
3445 async fn update_run_sets_output() {
3446 let store = InMemoryStore::new();
3447 let run = store
3448 .create_run(new_run_req("test"))
3449 .await
3450 .unwrap()
3451 .into_run();
3452
3453 store
3454 .update_run(
3455 run.id,
3456 RunUpdate {
3457 output: Some(json!({"verdict": "approved"})),
3458 ..RunUpdate::default()
3459 },
3460 )
3461 .await
3462 .unwrap();
3463
3464 let fetched = store.get_run(run.id).await.unwrap().unwrap();
3465 assert_eq!(fetched.output, Some(json!({"verdict": "approved"})));
3466 }
3467
3468 #[tokio::test]
3469 async fn update_run_without_output_keeps_previous_output() {
3470 let store = InMemoryStore::new();
3471 let run = store
3472 .create_run(new_run_req("test"))
3473 .await
3474 .unwrap()
3475 .into_run();
3476
3477 store
3478 .update_run(
3479 run.id,
3480 RunUpdate {
3481 output: Some(json!({"verdict": "approved"})),
3482 ..RunUpdate::default()
3483 },
3484 )
3485 .await
3486 .unwrap();
3487 store
3488 .update_run(
3489 run.id,
3490 RunUpdate {
3491 error: Some("boom".to_string()),
3492 ..RunUpdate::default()
3493 },
3494 )
3495 .await
3496 .unwrap();
3497
3498 let fetched = store.get_run(run.id).await.unwrap().unwrap();
3499 assert_eq!(fetched.output, Some(json!({"verdict": "approved"})));
3500 assert_eq!(fetched.error.as_deref(), Some("boom"));
3501 }
3502
3503 #[tokio::test]
3504 async fn update_run_output_last_write_wins() {
3505 let store = InMemoryStore::new();
3506 let run = store
3507 .create_run(new_run_req("test"))
3508 .await
3509 .unwrap()
3510 .into_run();
3511
3512 for verdict in ["first", "second"] {
3513 store
3514 .update_run(
3515 run.id,
3516 RunUpdate {
3517 output: Some(json!({ "verdict": verdict })),
3518 ..RunUpdate::default()
3519 },
3520 )
3521 .await
3522 .unwrap();
3523 }
3524
3525 let fetched = store.get_run(run.id).await.unwrap().unwrap();
3526 assert_eq!(fetched.output, Some(json!({"verdict": "second"})));
3527 }
3528
3529 async fn seed_user(store: &InMemoryStore, username: &str) -> Uuid {
3532 store
3533 .create_user(NewUser {
3534 email: format!("{username}@example.com"),
3535 username: username.to_string(),
3536 password_hash: "hash".to_string(),
3537 is_admin: Some(false),
3538 })
3539 .await
3540 .unwrap()
3541 .id
3542 }
3543
3544 async fn seed_api_key(store: &InMemoryStore, user_id: Uuid, name: &str) -> Uuid {
3545 store
3546 .create_api_key(NewApiKey {
3547 user_id,
3548 name: name.to_string(),
3549 key_hash: "hash".to_string(),
3550 key_prefix: "irfl_0000".to_string(),
3551 scopes: vec![ApiKeyScope::RunsWrite],
3552 expires_at: None,
3553 rate_limit_override: None,
3554 })
3555 .await
3556 .unwrap()
3557 .id
3558 }
3559
3560 fn run_req_by(actor: RunActor) -> NewRun {
3561 NewRun {
3562 created_by: Some(actor),
3563 ..new_run_req("test")
3564 }
3565 }
3566
3567 #[tokio::test]
3568 async fn create_run_without_actor_has_no_author() {
3569 let store = InMemoryStore::new();
3570 let run = store
3571 .create_run(new_run_req("test"))
3572 .await
3573 .unwrap()
3574 .into_run();
3575
3576 assert!(run.created_by.is_none());
3577 assert!(run.created_by_label.is_none());
3578 }
3579
3580 #[tokio::test]
3581 async fn create_run_by_user_resolves_username_as_label() {
3582 let store = InMemoryStore::new();
3583 let user_id = seed_user(&store, "alice").await;
3584
3585 let run = store
3586 .create_run(run_req_by(RunActor::User { user_id }))
3587 .await
3588 .unwrap()
3589 .into_run();
3590
3591 assert_eq!(run.created_by, Some(RunActor::User { user_id }));
3592 assert_eq!(run.created_by_label.as_deref(), Some("alice"));
3593 }
3594
3595 #[tokio::test]
3596 async fn create_run_by_api_key_resolves_key_and_owner_as_label() {
3597 let store = InMemoryStore::new();
3598 let user_id = seed_user(&store, "alice").await;
3599 let api_key_id = seed_api_key(&store, user_id, "ci-deploy").await;
3600
3601 let run = store
3602 .create_run(run_req_by(RunActor::ApiKey {
3603 api_key_id,
3604 user_id,
3605 }))
3606 .await
3607 .unwrap()
3608 .into_run();
3609
3610 assert_eq!(run.created_by_label.as_deref(), Some("ci-deploy (alice)"));
3611 }
3612
3613 #[tokio::test]
3614 async fn label_follows_api_key_rename() {
3615 let store = InMemoryStore::new();
3616 let user_id = seed_user(&store, "alice").await;
3617 let api_key_id = seed_api_key(&store, user_id, "ci-deploy").await;
3618 let run = store
3619 .create_run(run_req_by(RunActor::ApiKey {
3620 api_key_id,
3621 user_id,
3622 }))
3623 .await
3624 .unwrap()
3625 .into_run();
3626
3627 store
3628 .update_api_key(
3629 api_key_id,
3630 ApiKeyUpdate {
3631 name: Some("ci-release".to_string()),
3632 ..ApiKeyUpdate::default()
3633 },
3634 )
3635 .await
3636 .unwrap();
3637
3638 let reread = store.get_run(run.id).await.unwrap().unwrap();
3639 assert_eq!(
3640 reread.created_by_label.as_deref(),
3641 Some("ci-release (alice)")
3642 );
3643 }
3644
3645 #[tokio::test]
3646 async fn label_is_none_when_user_is_unknown() {
3647 let store = InMemoryStore::new();
3648 let run = store
3649 .create_run(run_req_by(RunActor::User {
3650 user_id: Uuid::now_v7(),
3651 }))
3652 .await
3653 .unwrap()
3654 .into_run();
3655
3656 assert!(run.created_by.is_some());
3657 assert!(run.created_by_label.is_none());
3658 }
3659
3660 #[tokio::test]
3661 async fn label_is_key_name_only_when_owner_is_unknown() {
3662 let store = InMemoryStore::new();
3663 let owner = seed_user(&store, "alice").await;
3664 let api_key_id = seed_api_key(&store, owner, "ci-deploy").await;
3665
3666 let run = store
3668 .create_run(run_req_by(RunActor::ApiKey {
3669 api_key_id,
3670 user_id: Uuid::now_v7(),
3671 }))
3672 .await
3673 .unwrap()
3674 .into_run();
3675
3676 assert_eq!(run.created_by_label.as_deref(), Some("ci-deploy"));
3677 }
3678
3679 #[tokio::test]
3680 async fn list_runs_filters_by_author() {
3681 let store = InMemoryStore::new();
3682 let alice = seed_user(&store, "alice").await;
3683 let bob = seed_user(&store, "bob").await;
3684
3685 store
3686 .create_run(run_req_by(RunActor::User { user_id: alice }))
3687 .await
3688 .unwrap()
3689 .into_run();
3690 store
3691 .create_run(run_req_by(RunActor::User { user_id: bob }))
3692 .await
3693 .unwrap()
3694 .into_run();
3695 store.create_run(new_run_req("anonymous")).await.unwrap();
3696
3697 let page = store
3698 .list_runs(
3699 RunFilter {
3700 created_by_user_id: Some(alice),
3701 ..RunFilter::default()
3702 },
3703 1,
3704 20,
3705 )
3706 .await
3707 .unwrap();
3708
3709 assert_eq!(page.total, 1);
3710 assert_eq!(page.items[0].created_by_label.as_deref(), Some("alice"));
3711 }
3712
3713 #[tokio::test]
3714 async fn list_runs_author_filter_matches_runs_from_the_users_api_keys() {
3715 let store = InMemoryStore::new();
3716 let alice = seed_user(&store, "alice").await;
3717 let api_key_id = seed_api_key(&store, alice, "ci-deploy").await;
3718
3719 store
3720 .create_run(run_req_by(RunActor::ApiKey {
3721 api_key_id,
3722 user_id: alice,
3723 }))
3724 .await
3725 .unwrap()
3726 .into_run();
3727
3728 let page = store
3729 .list_runs(
3730 RunFilter {
3731 created_by_user_id: Some(alice),
3732 ..RunFilter::default()
3733 },
3734 1,
3735 20,
3736 )
3737 .await
3738 .unwrap();
3739
3740 assert_eq!(page.total, 1);
3741 }
3742
3743 #[tokio::test]
3744 async fn list_runs_author_filter_excludes_unrelated_users() {
3745 let store = InMemoryStore::new();
3746 let alice = seed_user(&store, "alice").await;
3747
3748 store
3749 .create_run(run_req_by(RunActor::User { user_id: alice }))
3750 .await
3751 .unwrap()
3752 .into_run();
3753
3754 let page = store
3755 .list_runs(
3756 RunFilter {
3757 created_by_user_id: Some(Uuid::now_v7()),
3758 ..RunFilter::default()
3759 },
3760 1,
3761 20,
3762 )
3763 .await
3764 .unwrap();
3765
3766 assert_eq!(page.total, 0);
3767 }
3768
3769 #[tokio::test]
3770 async fn list_runs_without_author_filter_returns_every_run() {
3771 let store = InMemoryStore::new();
3772 let alice = seed_user(&store, "alice").await;
3773
3774 store
3775 .create_run(run_req_by(RunActor::User { user_id: alice }))
3776 .await
3777 .unwrap()
3778 .into_run();
3779 store.create_run(new_run_req("anonymous")).await.unwrap();
3780
3781 let page = store.list_runs(RunFilter::default(), 1, 20).await.unwrap();
3782 assert_eq!(page.total, 2);
3783 }
3784
3785 #[tokio::test]
3786 async fn pick_next_pending_resolves_author_label() {
3787 let store = InMemoryStore::new();
3788 let user_id = seed_user(&store, "alice").await;
3789 store
3790 .create_run(run_req_by(RunActor::User { user_id }))
3791 .await
3792 .unwrap()
3793 .into_run();
3794
3795 let picked = store.pick_next_pending(None).await.unwrap().unwrap();
3796 assert_eq!(picked.created_by_label.as_deref(), Some("alice"));
3797 }
3798
3799 #[tokio::test]
3802 async fn list_purgeable_runs_returns_old_terminal_runs() {
3803 let store = InMemoryStore::new();
3804 let old = create_terminal_run(&store, "deploy", RunStatus::Completed).await;
3805 store
3806 .set_run_created_at(old.id, Utc::now() - chrono::Duration::days(100))
3807 .await;
3808
3809 let policy = PurgePolicy {
3810 max_age_days: 90,
3811 max_runs_per_workflow: 10000,
3812 dry_run: false,
3813 };
3814 let result = store.list_purgeable_runs(&policy, 100).await.unwrap();
3815
3816 assert_eq!(result.len(), 1);
3817 assert_eq!(result[0].run_id, old.id);
3818 assert_eq!(result[0].reason, PurgeReason::TooOld);
3819 }
3820
3821 #[tokio::test]
3822 async fn list_purgeable_runs_ignores_non_terminal_states() {
3823 let store = InMemoryStore::new();
3824
3825 let pending = store
3827 .create_run(new_run_req("deploy"))
3828 .await
3829 .unwrap()
3830 .into_run();
3831 store
3832 .set_run_created_at(pending.id, Utc::now() - chrono::Duration::days(200))
3833 .await;
3834
3835 let running = store
3837 .create_run(new_run_req("deploy"))
3838 .await
3839 .unwrap()
3840 .into_run();
3841 store
3842 .update_run_status(running.id, RunStatus::Running)
3843 .await
3844 .unwrap();
3845 store
3846 .set_run_created_at(running.id, Utc::now() - chrono::Duration::days(200))
3847 .await;
3848
3849 let policy = PurgePolicy {
3850 max_age_days: 90,
3851 max_runs_per_workflow: 1,
3852 dry_run: false,
3853 };
3854 let result = store.list_purgeable_runs(&policy, 100).await.unwrap();
3855 assert!(result.is_empty());
3856 }
3857
3858 #[tokio::test]
3859 async fn list_purgeable_runs_returns_excess_per_workflow() {
3860 let store = InMemoryStore::new();
3861 let r1 = create_terminal_run(&store, "deploy", RunStatus::Completed).await;
3862 store
3863 .set_run_created_at(r1.id, Utc::now() - chrono::Duration::days(10))
3864 .await;
3865 let r2 = create_terminal_run(&store, "deploy", RunStatus::Completed).await;
3866 store
3867 .set_run_created_at(r2.id, Utc::now() - chrono::Duration::days(5))
3868 .await;
3869 let _r3 = create_terminal_run(&store, "deploy", RunStatus::Completed).await;
3870
3871 let policy = PurgePolicy {
3872 max_age_days: 365,
3873 max_runs_per_workflow: 2,
3874 dry_run: false,
3875 };
3876 let result = store.list_purgeable_runs(&policy, 100).await.unwrap();
3877
3878 assert_eq!(result.len(), 1);
3879 assert_eq!(result[0].run_id, r1.id);
3880 assert_eq!(result[0].reason, PurgeReason::ExceedsWorkflowLimit);
3881 }
3882
3883 #[tokio::test]
3886 async fn delete_run_removes_run_and_associated_data() {
3887 use crate::artifact_store::ArtifactStore;
3888 use crate::entities::{NewStep, StepKind};
3889
3890 let store = InMemoryStore::new();
3891 let run = create_terminal_run(&store, "deploy", RunStatus::Completed).await;
3892 let step = store
3893 .create_step(NewStep {
3894 run_id: run.id,
3895 trace_id: step_trace_id(run.id, "build", 0),
3896 name: "build".to_string(),
3897 kind: StepKind::Shell,
3898 position: 0,
3899 input: None,
3900 is_error_handler: false,
3901 })
3902 .await
3903 .unwrap();
3904
3905 let artifact_id = Uuid::now_v7();
3906 store
3907 .create_artifact(crate::entities::NewArtifact {
3908 id: artifact_id,
3909 run_id: run.id,
3910 step_id: step.id,
3911 name: "report.html".to_string(),
3912 storage_key: format!("artifacts/{}/{}/{}", run.id, step.id, artifact_id),
3913 content_type: "text/html".to_string(),
3914 size_bytes: 42,
3915 sha256: "0".repeat(64),
3916 })
3917 .await
3918 .unwrap();
3919
3920 let keys = store.delete_run(run.id).await.unwrap();
3921
3922 assert_eq!(keys.len(), 1);
3923 assert!(keys[0].contains(&artifact_id.to_string()));
3924 assert!(store.get_run(run.id).await.unwrap().is_none());
3925 assert!(store.list_steps(run.id).await.unwrap().is_empty());
3926 assert!(
3927 store
3928 .list_artifacts_for_run(run.id)
3929 .await
3930 .unwrap()
3931 .is_empty()
3932 );
3933 }
3934
3935 #[tokio::test]
3936 async fn delete_run_not_found() {
3937 let store = InMemoryStore::new();
3938 let err = store.delete_run(Uuid::now_v7()).await.unwrap_err();
3939 assert!(matches!(err, StoreError::RunNotFound(_)));
3940 }
3941
3942 fn vote(user_id: Uuid, name: &str) -> StepApproval {
3945 StepApproval {
3946 user_id,
3947 approved_by: name.to_string(),
3948 at: Utc::now(),
3949 }
3950 }
3951
3952 #[tokio::test]
3953 async fn record_step_approval_appends_distinct_voters() {
3954 let store = InMemoryStore::new();
3955 let run = store
3956 .create_run(new_run_req("test"))
3957 .await
3958 .unwrap()
3959 .into_run();
3960 let step = store
3961 .create_step(new_step_req(run.id, "gate", 0))
3962 .await
3963 .unwrap();
3964 assert!(step.approvals.is_empty());
3965 assert!(step.approval_requirement.is_none());
3966
3967 let alice = Uuid::now_v7();
3968 let bob = Uuid::now_v7();
3969 let after_first = store
3970 .record_step_approval(step.id, vote(alice, "alice"))
3971 .await
3972 .unwrap();
3973 assert_eq!(after_first.approvals.len(), 1);
3974
3975 let after_second = store
3976 .record_step_approval(step.id, vote(bob, "bob"))
3977 .await
3978 .unwrap();
3979 assert_eq!(after_second.approvals.len(), 2);
3980 assert_eq!(after_second.approvals[0].user_id, alice);
3981 assert_eq!(after_second.approvals[1].user_id, bob);
3982 }
3983
3984 #[tokio::test]
3985 async fn record_step_approval_ignores_same_user() {
3986 let store = InMemoryStore::new();
3987 let run = store
3988 .create_run(new_run_req("test"))
3989 .await
3990 .unwrap()
3991 .into_run();
3992 let step = store
3993 .create_step(new_step_req(run.id, "gate", 0))
3994 .await
3995 .unwrap();
3996
3997 let alice = Uuid::now_v7();
3998 store
3999 .record_step_approval(step.id, vote(alice, "alice"))
4000 .await
4001 .unwrap();
4002 let again = store
4003 .record_step_approval(step.id, vote(alice, "alice-key"))
4004 .await
4005 .unwrap();
4006
4007 assert_eq!(again.approvals.len(), 1);
4008 assert_eq!(again.approvals[0].approved_by, "alice");
4009 }
4010
4011 #[tokio::test]
4012 async fn record_step_approval_unknown_step_is_not_found() {
4013 let store = InMemoryStore::new();
4014 let err = store
4015 .record_step_approval(Uuid::now_v7(), vote(Uuid::now_v7(), "alice"))
4016 .await
4017 .unwrap_err();
4018 assert!(matches!(err, StoreError::StepNotFound(_)));
4019 }
4020
4021 #[tokio::test]
4022 async fn update_step_sets_approval_requirement() {
4023 let store = InMemoryStore::new();
4024 let run = store
4025 .create_run(new_run_req("test"))
4026 .await
4027 .unwrap()
4028 .into_run();
4029 let step = store
4030 .create_step(new_step_req(run.id, "gate", 0))
4031 .await
4032 .unwrap();
4033 let requirement = ApprovalRequirement {
4034 required_approvers: 3,
4035 ..ApprovalRequirement::default()
4036 };
4037
4038 store
4039 .update_step(
4040 step.id,
4041 StepUpdate {
4042 approval_requirement: Some(requirement.clone()),
4043 ..StepUpdate::default()
4044 },
4045 )
4046 .await
4047 .unwrap();
4048
4049 let fetched = store.get_step(step.id).await.unwrap().unwrap();
4050 assert_eq!(fetched.approval_requirement, Some(requirement));
4051 }
4052
4053 async fn sleeping_run(store: &InMemoryStore, scheduled_at: DateTime<Utc>) -> Run {
4054 let run = store
4055 .create_run(new_run_req("sleepy"))
4056 .await
4057 .unwrap()
4058 .into_run();
4059 store
4060 .update_run_status(run.id, RunStatus::Running)
4061 .await
4062 .unwrap();
4063 store
4064 .update_run(
4065 run.id,
4066 RunUpdate {
4067 status: Some(RunStatus::Sleeping),
4068 scheduled_at: Some(scheduled_at),
4069 ..RunUpdate::default()
4070 },
4071 )
4072 .await
4073 .unwrap();
4074 store.get_run(run.id).await.unwrap().unwrap()
4075 }
4076
4077 #[tokio::test]
4078 async fn claim_due_sleeping_runs_requeues_due_runs() {
4079 let store = InMemoryStore::new();
4080 let due = sleeping_run(&store, Utc::now() - TimeDelta::seconds(5)).await;
4081
4082 let woken = store.claim_due_sleeping_runs(10).await.unwrap();
4083 assert_eq!(woken.len(), 1);
4084 assert_eq!(woken[0].id, due.id);
4085 assert_eq!(woken[0].status.state, RunStatus::Pending);
4086 assert!(woken[0].scheduled_at.is_none());
4087
4088 let fetched = store.get_run(due.id).await.unwrap().unwrap();
4089 assert_eq!(fetched.status.state, RunStatus::Pending);
4090 assert!(fetched.scheduled_at.is_none());
4091
4092 assert!(store.claim_due_sleeping_runs(10).await.unwrap().is_empty());
4094 }
4095
4096 #[tokio::test]
4097 async fn claim_due_sleeping_runs_skips_future_runs() {
4098 let store = InMemoryStore::new();
4099 let future = sleeping_run(&store, Utc::now() + TimeDelta::hours(1)).await;
4100
4101 assert!(store.claim_due_sleeping_runs(10).await.unwrap().is_empty());
4102 let fetched = store.get_run(future.id).await.unwrap().unwrap();
4103 assert_eq!(fetched.status.state, RunStatus::Sleeping);
4104 }
4105
4106 #[tokio::test]
4107 async fn claim_due_sleeping_runs_honours_limit_oldest_first() {
4108 let store = InMemoryStore::new();
4109 let older = sleeping_run(&store, Utc::now() - TimeDelta::seconds(20)).await;
4110 let newer = sleeping_run(&store, Utc::now() - TimeDelta::seconds(10)).await;
4111
4112 let woken = store.claim_due_sleeping_runs(1).await.unwrap();
4113 assert_eq!(woken.len(), 1);
4114 assert_eq!(woken[0].id, older.id);
4115
4116 let woken = store.claim_due_sleeping_runs(1).await.unwrap();
4117 assert_eq!(woken[0].id, newer.id);
4118 }
4119}