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