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