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