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