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