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