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