Skip to main content

ironflow_store/memory/
run_store.rs

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