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