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