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