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