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