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