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