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