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