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