Skip to main content

ironflow_store/memory/
run_store.rs

1use std::collections::{BTreeMap, HashMap, HashSet};
2
3use chrono::{DateTime, Duration, Utc};
4use rust_decimal::Decimal;
5use uuid::Uuid;
6
7use crate::entities::{
8    ApiKey, ConcurrencyGroupBacklog, IDEMPOTENCY_WINDOW, LeaseRequest, LeaseUpdate, NewRun,
9    NewStep, NewStepDependency, Page, PurgePolicy, PurgeReason, PurgeableRun, ReapedRun, Run,
10    RunActor, RunCreation, RunFilter, RunStats, RunStatus, RunUpdate, StatsHistoryBucket,
11    StatsHistoryFilter, Step, StepApproval, StepDependency, StepStatus, StepUpdate, TriggerKind,
12    User, WorkerCapabilities, normalize_worker_tags, validate_concurrency_limits,
13    validate_worker_tags,
14};
15use crate::error::StoreError;
16use crate::store::{LEASE_EXPIRED_ERROR, RunStore, StoreFuture};
17
18use super::descendants::active_descendants;
19use super::stats_history::aggregate_history_buckets;
20use super::{InMemoryStore, State};
21
22/// Resolve [`Run::created_by_label`] from the current users and API keys.
23///
24/// Mirrors the `LEFT JOIN` the PostgreSQL store performs at read time, so a
25/// renamed API key or a renamed user is reflected immediately.
26fn resolve_created_by_label(
27    actor: Option<&RunActor>,
28    users: &HashMap<Uuid, User>,
29    api_keys: &HashMap<Uuid, ApiKey>,
30) -> Option<String> {
31    let actor = actor?;
32    let username = users.get(&actor.user_id()).map(|u| u.username.clone());
33
34    match actor.api_key_id() {
35        None => username,
36        Some(api_key_id) => {
37            let key_name = api_keys.get(&api_key_id).map(|k| k.name.clone())?;
38            Some(match username {
39                Some(username) => format!("{key_name} ({username})"),
40                None => key_name,
41            })
42        }
43    }
44}
45
46/// Return a clone of `run` with its display label resolved against `state`.
47fn run_with_label(run: &Run, state: &State) -> Run {
48    let mut run = run.clone();
49    run.created_by_label =
50        resolve_created_by_label(run.created_by.as_ref(), &state.users, &state.api_keys);
51    run
52}
53
54/// Drop the worker lease held on a run.
55fn clear_lease(run: &mut Run) {
56    run.worker_id = None;
57    run.lease_expires_at = None;
58}
59
60/// Number of root runs in `Running` that carry `group`.
61///
62/// Sub-workflow runs execute inside their parent's slot and are never counted.
63fn running_count(runs: &HashMap<Uuid, Run>, group: &str) -> u64 {
64    runs.values()
65        .filter(|r| {
66            r.status.state == RunStatus::Running
67                && !matches!(r.trigger, TriggerKind::Workflow)
68                && r.concurrency_limits.iter().any(|l| l.group == group)
69        })
70        .count() as u64
71}
72
73/// Whether at least one of the run's concurrency groups is saturated for the
74/// run's own limit.
75fn is_blocked(run: &Run, runs: &HashMap<Uuid, Run>) -> bool {
76    run.concurrency_limits
77        .iter()
78        .any(|l| running_count(runs, &l.group) >= u64::from(l.limit))
79}
80
81/// Whether a run waits for a pick: pending or retrying, and due.
82fn is_due(run: &Run, now: DateTime<Utc>) -> bool {
83    matches!(run.status.state, RunStatus::Pending | RunStatus::Retrying)
84        && run.scheduled_at.is_none_or(|at| at <= now)
85}
86
87fn run_matches_filter(run: &Run, filter: &RunFilter, steps: &HashMap<Uuid, Step>) -> bool {
88    if let Some(ref wf) = filter.workflow_name
89        && !run
90            .workflow_name
91            .to_lowercase()
92            .contains(&wf.to_lowercase())
93    {
94        return false;
95    }
96    if let Some(ref status) = filter.status
97        && &run.status.state != status
98    {
99        return false;
100    }
101    if let Some(after) = filter.created_after
102        && run.created_at < after
103    {
104        return false;
105    }
106    if let Some(before) = filter.created_before
107        && run.created_at > before
108    {
109        return false;
110    }
111    if let Some(has_steps) = filter.has_steps
112        && matches!(
113            run.status.state,
114            RunStatus::Completed | RunStatus::Cancelled
115        )
116    {
117        let run_has_steps = steps.values().any(|s| s.run_id == run.id);
118        if has_steps != run_has_steps {
119            return false;
120        }
121    }
122    if let Some(ref labels) = filter.labels {
123        for (key, value) in labels {
124            if run.labels.get(key) != Some(value) {
125                return false;
126            }
127        }
128    }
129    if let Some(user_id) = filter.created_by_user_id
130        && run.created_by.as_ref().map(RunActor::user_id) != Some(user_id)
131    {
132        return false;
133    }
134    if let Some(ref group) = filter.concurrency_group
135        && !run.concurrency_limits.iter().any(|l| &l.group == group)
136    {
137        return false;
138    }
139    if let Some(ref caps) = filter.eligible_for
140        && !caps.can_take(&run.workflow_name, &run.worker_tags)
141    {
142        return false;
143    }
144    true
145}
146
147/// Insert a run into the locked state, honouring its idempotency and
148/// concurrency keys. Writes nothing when it returns an error.
149///
150/// Shared by [`RunStore::create_run`] and the schedule firing, which must
151/// create the run under the same lock as the schedule update.
152pub(super) fn insert_run(state: &mut State, req: NewRun) -> Result<RunCreation, StoreError> {
153    validate_concurrency_limits(&req.concurrency_limits)?;
154    validate_worker_tags(&req.worker_tags)?;
155    let now = Utc::now();
156
157    if let Some(ref key) = req.idempotency_key
158        && let Some(existing) = state
159            .idempotency_keys
160            .get(key)
161            .and_then(|id| state.runs.get(id))
162    {
163        if now - existing.created_at < IDEMPOTENCY_WINDOW {
164            return Ok(RunCreation::Existing(run_with_label(existing, state)));
165        }
166        // The key outlived its window: release it from the stale run.
167        let stale_id = existing.id;
168        state.idempotency_keys.remove(key);
169        if let Some(stale) = state.runs.get_mut(&stale_id) {
170            stale.idempotency_key = None;
171        }
172    }
173
174    if let Some(ref key) = req.concurrency_key
175        && let Some(holder) = state
176            .runs
177            .values()
178            .filter(|r| {
179                r.concurrency_key.as_deref() == Some(key.as_str()) && !r.status.state.is_terminal()
180            })
181            .min_by_key(|r| r.created_at)
182    {
183        return Err(StoreError::ConcurrencyConflict {
184            key: key.clone(),
185            run_id: holder.id,
186        });
187    }
188
189    let run = Run {
190        id: Uuid::now_v7(),
191        workflow_name: req.workflow_name,
192        status: crate::entities::FsmState::new(RunStatus::Pending, Uuid::now_v7()),
193        trigger: req.trigger,
194        payload: req.payload,
195        error: None,
196        retry_count: 0,
197        max_retries: req.max_retries,
198        cost_usd: Decimal::ZERO,
199        duration_ms: 0,
200        created_at: now,
201        updated_at: now,
202        started_at: None,
203        completed_at: None,
204        handler_version: req.handler_version,
205        labels: req.labels,
206        scheduled_at: req.scheduled_at,
207        created_by: req.created_by,
208        created_by_label: None,
209        idempotency_key: req.idempotency_key.clone(),
210        concurrency_key: req.concurrency_key,
211        concurrency_limits: req.concurrency_limits,
212        max_cost_usd: req.max_cost_usd,
213        worker_id: None,
214        lease_expires_at: None,
215        output: None,
216        lease_recoveries: 0,
217        capacity_wait_kind: None,
218        worker_tags: normalize_worker_tags(req.worker_tags),
219    };
220
221    if let Some(key) = req.idempotency_key {
222        state.idempotency_keys.insert(key, run.id);
223    }
224    state.runs.insert(run.id, run.clone());
225    Ok(RunCreation::Created(run_with_label(&run, state)))
226}
227
228impl RunStore for InMemoryStore {
229    fn create_run(&self, req: NewRun) -> StoreFuture<'_, RunCreation> {
230        Box::pin(async move {
231            // Single critical section: the key lookup and the insert cannot be
232            // interleaved by a concurrent call sharing the same key.
233            let mut state = self.state.write().await;
234            insert_run(&mut state, req)
235        })
236    }
237
238    fn find_run_by_idempotency_key(&self, key: &str) -> StoreFuture<'_, Option<Run>> {
239        let key = key.to_string();
240        Box::pin(async move {
241            let now = Utc::now();
242            let state = self.state.read().await;
243            Ok(state
244                .idempotency_keys
245                .get(&key)
246                .and_then(|id| state.runs.get(id))
247                .filter(|run| now - run.created_at < IDEMPOTENCY_WINDOW)
248                .map(|run| run_with_label(run, &state)))
249        })
250    }
251
252    fn get_run(&self, id: Uuid) -> StoreFuture<'_, Option<Run>> {
253        Box::pin(async move {
254            let state = self.state.read().await;
255            Ok(state.runs.get(&id).map(|r| run_with_label(r, &state)))
256        })
257    }
258
259    fn list_runs(&self, filter: RunFilter, page: u32, per_page: u32) -> StoreFuture<'_, Page<Run>> {
260        Box::pin(async move {
261            let state = self.state.read().await;
262
263            let mut runs: Vec<&Run> = state
264                .runs
265                .values()
266                .filter(|r| run_matches_filter(r, &filter, &state.steps))
267                .collect();
268
269            // Sort newest first.
270            runs.sort_by_key(|r| std::cmp::Reverse(r.created_at));
271
272            let total = runs.len() as u64;
273            let page = page.max(1);
274            let per_page = per_page.clamp(1, 100);
275            let offset = ((page - 1) * per_page) as usize;
276            let items: Vec<Run> = runs
277                .into_iter()
278                .skip(offset)
279                .take(per_page as usize)
280                .map(|r| run_with_label(r, &state))
281                .collect();
282
283            Ok(Page {
284                items,
285                total,
286                page,
287                per_page,
288            })
289        })
290    }
291
292    fn update_run_status(&self, id: Uuid, new_status: RunStatus) -> StoreFuture<'_, ()> {
293        Box::pin(async move {
294            let mut state = self.state.write().await;
295            let run = state.runs.get_mut(&id).ok_or(StoreError::RunNotFound(id))?;
296
297            if !run.status.state.can_transition_to(&new_status) {
298                return Err(StoreError::InvalidTransition {
299                    from: run.status.state,
300                    to: new_status,
301                });
302            }
303
304            if run.status.state == new_status && new_status.is_terminal() {
305                return Ok(());
306            }
307
308            let now = Utc::now();
309            run.status.state = new_status;
310            run.updated_at = now;
311
312            if new_status == RunStatus::Running && run.started_at.is_none() {
313                run.started_at = Some(now);
314            }
315            if new_status.is_terminal() {
316                run.completed_at = Some(now);
317            }
318            if new_status != RunStatus::Running {
319                clear_lease(run);
320            }
321            if new_status != RunStatus::Sleeping {
322                run.capacity_wait_kind = None;
323            }
324
325            Ok(())
326        })
327    }
328
329    fn update_run(&self, id: Uuid, update: RunUpdate) -> StoreFuture<'_, ()> {
330        Box::pin(async move {
331            let mut state = self.state.write().await;
332            let run = state.runs.get_mut(&id).ok_or(StoreError::RunNotFound(id))?;
333
334            let now = Utc::now();
335
336            if let Some(status) = update.status {
337                if !run.status.state.can_transition_to(&status) {
338                    return Err(StoreError::InvalidTransition {
339                        from: run.status.state,
340                        to: status,
341                    });
342                }
343                if !(run.status.state == status && status.is_terminal()) {
344                    run.status.state = status;
345                    if status == RunStatus::Running && run.started_at.is_none() {
346                        run.started_at = Some(now);
347                    }
348                    if status.is_terminal() {
349                        run.completed_at = Some(now);
350                    }
351                    if status != RunStatus::Running {
352                        clear_lease(run);
353                    }
354                }
355                // The kind only means something while the run sleeps on capacity.
356                run.capacity_wait_kind = if status == RunStatus::Sleeping {
357                    update.capacity_wait_kind.clone()
358                } else {
359                    None
360                };
361            }
362
363            // After the status block, so `Running` + `Set` ends with a lease.
364            match update.lease {
365                Some(LeaseUpdate::Set {
366                    worker_id,
367                    expires_at,
368                }) => {
369                    run.worker_id = Some(worker_id);
370                    run.lease_expires_at = Some(expires_at);
371                }
372                Some(LeaseUpdate::Release) => clear_lease(run),
373                None => {}
374            }
375
376            if let Some(error) = update.error {
377                run.error = Some(error);
378            }
379            if update.increment_retry {
380                run.retry_count += 1;
381            }
382            if let Some(cost) = update.cost_usd {
383                run.cost_usd = cost;
384            }
385            if let Some(dur) = update.duration_ms {
386                run.duration_ms = dur;
387            }
388            if let Some(started) = update.started_at {
389                run.started_at = Some(started);
390            }
391            if let Some(completed) = update.completed_at {
392                run.completed_at = Some(completed);
393            }
394            if let Some(scheduled) = update.scheduled_at {
395                run.scheduled_at = Some(scheduled);
396            }
397            if let Some(output) = update.output {
398                run.output = Some(output);
399            }
400
401            run.updated_at = now;
402            Ok(())
403        })
404    }
405
406    fn list_active_descendants(&self, run_id: Uuid) -> StoreFuture<'_, Vec<Run>> {
407        Box::pin(async move {
408            let state = self.state.read().await;
409            Ok(active_descendants(&state.runs, run_id))
410        })
411    }
412
413    fn pick_next_pending_for(
414        &self,
415        lease: Option<LeaseRequest>,
416        capabilities: Option<WorkerCapabilities>,
417    ) -> StoreFuture<'_, Option<Run>> {
418        Box::pin(async move {
419            let mut state = self.state.write().await;
420            let now = Utc::now();
421
422            // Find the oldest run waiting for execution whose scheduled_at has
423            // passed (or is None). `Retrying` runs are runs whose automatic retry
424            // backoff has been armed: they become eligible again once
425            // `scheduled_at` has passed. Runs held back by a saturated
426            // concurrency group are skipped, so they never block younger runs.
427            // The group check and the transition below happen under the same
428            // write lock, so concurrent pickers cannot overshoot a limit.
429            // Runs the worker cannot take (unregistered workflow, missing tag)
430            // are skipped the same way.
431            let oldest_id = state
432                .runs
433                .values()
434                .filter(|r| is_due(r, now) && !is_blocked(r, &state.runs))
435                .filter(|r| {
436                    capabilities
437                        .as_ref()
438                        .is_none_or(|c| c.can_take(&r.workflow_name, &r.worker_tags))
439                })
440                .min_by_key(|r| r.created_at)
441                .map(|r| r.id);
442
443            let Some(id) = oldest_id else {
444                return Ok(None);
445            };
446
447            // Transition to Running and attach the lease atomically.
448            let run = state.runs.get_mut(&id).expect("run exists");
449            let now = Utc::now();
450            run.status.state = RunStatus::Running;
451            run.started_at = Some(now);
452            run.updated_at = now;
453            match lease {
454                Some(lease) => {
455                    run.lease_expires_at = Some(lease.expires_at(now));
456                    run.worker_id = Some(lease.worker_id);
457                }
458                None => clear_lease(run),
459            }
460            let run = run.clone();
461
462            Ok(Some(run_with_label(&run, &state)))
463        })
464    }
465
466    fn count_blocked_runs_by_group(&self) -> StoreFuture<'_, Vec<ConcurrencyGroupBacklog>> {
467        Box::pin(async move {
468            let state = self.state.read().await;
469            let now = Utc::now();
470
471            let mut blocked: BTreeMap<String, u64> = BTreeMap::new();
472            for run in state.runs.values().filter(|r| is_due(r, now)) {
473                for limit in &run.concurrency_limits {
474                    if running_count(&state.runs, &limit.group) >= u64::from(limit.limit) {
475                        *blocked.entry(limit.group.clone()).or_default() += 1;
476                    }
477                }
478            }
479
480            Ok(blocked
481                .into_iter()
482                .map(|(group, blocked_runs)| ConcurrencyGroupBacklog {
483                    group,
484                    blocked_runs,
485                })
486                .collect())
487        })
488    }
489
490    fn renew_lease(&self, id: Uuid, lease: LeaseRequest) -> StoreFuture<'_, DateTime<Utc>> {
491        Box::pin(async move {
492            let mut state = self.state.write().await;
493            let run = state.runs.get_mut(&id).ok_or(StoreError::RunNotFound(id))?;
494
495            if run.status.state != RunStatus::Running
496                || run.worker_id.as_deref() != Some(lease.worker_id.as_str())
497            {
498                return Err(StoreError::LeaseLost {
499                    run_id: id,
500                    held_by: run.worker_id.clone(),
501                });
502            }
503
504            let now = Utc::now();
505            let expires_at = lease.expires_at(now);
506            run.lease_expires_at = Some(expires_at);
507            run.updated_at = now;
508
509            Ok(expires_at)
510        })
511    }
512
513    fn reap_expired_leases(&self, limit: u32) -> StoreFuture<'_, Vec<ReapedRun>> {
514        Box::pin(async move {
515            let mut state = self.state.write().await;
516            let now = Utc::now();
517
518            let mut expired: Vec<Uuid> = state
519                .runs
520                .values()
521                .filter(|r| {
522                    r.status.state == RunStatus::Running
523                        && r.lease_expires_at.is_some_and(|at| at < now)
524                })
525                .map(|r| r.id)
526                .collect();
527            expired.sort_unstable();
528            expired.truncate(limit as usize);
529
530            let mut reaped = Vec::with_capacity(expired.len());
531            for id in expired {
532                let run = state.runs.get_mut(&id).expect("run exists");
533                run.lease_recoveries += 1;
534                clear_lease(run);
535                run.updated_at = now;
536
537                let to = if run.lease_recoveries > run.max_retries {
538                    run.status.state = RunStatus::Failed;
539                    run.error = Some(LEASE_EXPIRED_ERROR.to_string());
540                    run.completed_at = Some(now);
541                    RunStatus::Failed
542                } else {
543                    run.status.state = RunStatus::Pending;
544                    RunStatus::Pending
545                };
546
547                reaped.push(ReapedRun {
548                    run: run.clone(),
549                    from: RunStatus::Running,
550                    to,
551                });
552            }
553
554            Ok(reaped)
555        })
556    }
557
558    fn claim_due_approval_deadlines(&self, limit: u32) -> StoreFuture<'_, Vec<Step>> {
559        Box::pin(async move {
560            let mut state = self.state.write().await;
561            let now = Utc::now();
562
563            let mut due: Vec<(DateTime<Utc>, Uuid)> = state
564                .steps
565                .values()
566                .filter(|s| {
567                    s.status.state == StepStatus::AwaitingApproval
568                        && s.approval_deadline_at.is_some_and(|at| at <= now)
569                })
570                .map(|s| (s.approval_deadline_at.expect("deadline is set"), s.id))
571                .collect();
572            due.sort_unstable();
573            due.truncate(limit as usize);
574
575            let mut claimed = Vec::with_capacity(due.len());
576            for (_, id) in due {
577                let step = state.steps.get_mut(&id).expect("step exists");
578                // Clone before clearing so the caller still sees the deadline
579                // that fired.
580                claimed.push(step.clone());
581                step.approval_deadline_at = None;
582                step.updated_at = now;
583            }
584
585            Ok(claimed)
586        })
587    }
588
589    fn claim_due_sleeping_runs(&self, limit: u32) -> StoreFuture<'_, Vec<Run>> {
590        Box::pin(async move {
591            let mut state = self.state.write().await;
592            let now = Utc::now();
593
594            let mut due: Vec<(DateTime<Utc>, Uuid)> = state
595                .runs
596                .values()
597                .filter(|r| r.status.state == RunStatus::Sleeping)
598                .filter_map(|r| r.scheduled_at.filter(|at| *at <= now).map(|at| (at, r.id)))
599                .collect();
600            due.sort_unstable();
601            due.truncate(limit as usize);
602
603            let mut woken = Vec::with_capacity(due.len());
604            for (_, id) in due {
605                let run = state.runs.get_mut(&id).expect("run exists");
606                run.status.state = RunStatus::Pending;
607                run.scheduled_at = None;
608                run.capacity_wait_kind = None;
609                run.updated_at = now;
610                let run = run.clone();
611                woken.push(run_with_label(&run, &state));
612            }
613
614            Ok(woken)
615        })
616    }
617
618    fn list_purgeable_runs(
619        &self,
620        policy: &PurgePolicy,
621        batch_size: u32,
622    ) -> StoreFuture<'_, Vec<PurgeableRun>> {
623        let max_age_days = policy.max_age_days;
624        let max_runs_per_workflow = policy.max_runs_per_workflow;
625        Box::pin(async move {
626            let state = self.state.read().await;
627            let cutoff = Utc::now() - Duration::days(i64::from(max_age_days));
628            let mut result: Vec<PurgeableRun> = Vec::new();
629            let mut seen: HashSet<Uuid> = HashSet::new();
630
631            for run in state.runs.values() {
632                if run.status.state.is_terminal() && run.created_at < cutoff {
633                    seen.insert(run.id);
634                    result.push(PurgeableRun {
635                        run_id: run.id,
636                        workflow_name: run.workflow_name.clone(),
637                        reason: PurgeReason::TooOld,
638                    });
639                }
640            }
641
642            let mut by_workflow: HashMap<&str, Vec<&Run>> = HashMap::new();
643            for run in state.runs.values() {
644                if run.status.state.is_terminal() {
645                    by_workflow.entry(&run.workflow_name).or_default().push(run);
646                }
647            }
648            for (_, mut runs) in by_workflow {
649                if runs.len() > max_runs_per_workflow as usize {
650                    runs.sort_by_key(|r| r.created_at);
651                    let excess = runs.len() - max_runs_per_workflow as usize;
652                    for run in runs.into_iter().take(excess) {
653                        if seen.insert(run.id) {
654                            result.push(PurgeableRun {
655                                run_id: run.id,
656                                workflow_name: run.workflow_name.clone(),
657                                reason: PurgeReason::ExceedsWorkflowLimit,
658                            });
659                        }
660                    }
661                }
662            }
663
664            result.sort_by_key(|p| p.run_id);
665            result.truncate(batch_size as usize);
666            Ok(result)
667        })
668    }
669
670    fn delete_run(&self, id: Uuid) -> StoreFuture<'_, Vec<String>> {
671        Box::pin(async move {
672            let mut state = self.state.write().await;
673
674            if !state.runs.contains_key(&id) {
675                return Err(StoreError::RunNotFound(id));
676            }
677
678            // Collect artifact storage keys before removing them.
679            let storage_keys: Vec<String> = state
680                .artifacts
681                .values()
682                .filter(|a| a.run_id == id)
683                .map(|a| a.storage_key.clone())
684                .collect();
685
686            // Remove artifacts.
687            state.artifacts.retain(|_, a| a.run_id != id);
688
689            // Remove step dependencies (both sides).
690            let step_ids: Vec<Uuid> = state
691                .steps
692                .values()
693                .filter(|s| s.run_id == id)
694                .map(|s| s.id)
695                .collect();
696            state
697                .step_dependencies
698                .retain(|d| !step_ids.contains(&d.step_id) && !step_ids.contains(&d.depends_on));
699
700            // Remove steps.
701            state.steps.retain(|_, s| s.run_id != id);
702
703            // Remove idempotency key pointing to this run.
704            state.idempotency_keys.retain(|_, &mut run_id| run_id != id);
705
706            // Remove the run itself.
707            state.runs.remove(&id);
708
709            Ok(storage_keys)
710        })
711    }
712
713    fn create_step(&self, req: NewStep) -> StoreFuture<'_, Step> {
714        Box::pin(async move {
715            let mut state = self.state.write().await;
716
717            let attempt = state
718                .runs
719                .get(&req.run_id)
720                .ok_or(StoreError::RunNotFound(req.run_id))?
721                .retry_count
722                + 1;
723
724            let now = Utc::now();
725            let step = Step {
726                id: Uuid::now_v7(),
727                trace_id: req.trace_id,
728                run_id: req.run_id,
729                name: req.name,
730                kind: req.kind,
731                position: req.position,
732                status: crate::entities::FsmState::new(StepStatus::Pending, Uuid::now_v7()),
733                attempt,
734                input: req.input,
735                output: None,
736                error: None,
737                duration_ms: 0,
738                cost_usd: Decimal::ZERO,
739                input_tokens: None,
740                cache_read_input_tokens: None,
741                cache_creation_input_tokens: None,
742                output_tokens: None,
743                created_at: now,
744                updated_at: now,
745                started_at: None,
746                completed_at: None,
747                debug_messages: None,
748                is_error_handler: req.is_error_handler,
749                approval_deadline_at: None,
750                approval_stage: 0,
751                approval_assignee: None,
752                approval_requirement: None,
753                approvals: Vec::new(),
754                account_id: None,
755                environment_id: None,
756                session_id: None,
757            };
758
759            state.steps.insert(step.id, step.clone());
760            Ok(step)
761        })
762    }
763
764    fn update_step(&self, id: Uuid, update: StepUpdate) -> StoreFuture<'_, ()> {
765        Box::pin(async move {
766            let mut state = self.state.write().await;
767            let step = state
768                .steps
769                .get_mut(&id)
770                .ok_or(StoreError::StepNotFound(id))?;
771
772            let now = Utc::now();
773
774            if let Some(status) = update.status {
775                if !matches!(
776                    (step.status.state, status),
777                    (StepStatus::Pending, StepStatus::Running)
778                        | (StepStatus::Pending, StepStatus::Skipped)
779                        | (StepStatus::Running, StepStatus::Completed)
780                        | (StepStatus::Running, StepStatus::Failed)
781                        | (StepStatus::Running, StepStatus::AwaitingApproval)
782                        | (StepStatus::AwaitingApproval, StepStatus::Running)
783                        | (StepStatus::AwaitingApproval, StepStatus::Completed)
784                        | (StepStatus::AwaitingApproval, StepStatus::Failed)
785                        | (StepStatus::AwaitingApproval, StepStatus::Rejected)
786                ) {
787                    return Err(StoreError::Database(format!(
788                        "invalid step status transition: {:?} -> {:?}",
789                        step.status.state, status
790                    )));
791                }
792                step.status.state = status;
793            }
794            if let Some(output) = update.output {
795                step.output = Some(output);
796            }
797            if let Some(error) = update.error {
798                step.error = Some(error);
799            }
800            if let Some(dur) = update.duration_ms {
801                step.duration_ms = dur;
802            }
803            if let Some(cost) = update.cost_usd {
804                step.cost_usd = cost;
805            }
806            if let Some(tokens) = update.input_tokens {
807                step.input_tokens = Some(tokens);
808            }
809            if let Some(tokens) = update.cache_read_input_tokens {
810                step.cache_read_input_tokens = Some(tokens);
811            }
812            if let Some(tokens) = update.cache_creation_input_tokens {
813                step.cache_creation_input_tokens = Some(tokens);
814            }
815            if let Some(account_id) = update.account_id {
816                step.account_id = Some(account_id);
817            }
818            if let Some(environment_id) = update.environment_id {
819                step.environment_id = Some(environment_id);
820            }
821            if let Some(session_id) = update.session_id {
822                step.session_id = Some(session_id);
823            }
824            if let Some(tokens) = update.output_tokens {
825                step.output_tokens = Some(tokens);
826            }
827            if let Some(started) = update.started_at {
828                step.started_at = Some(started);
829            }
830            if let Some(completed) = update.completed_at {
831                step.completed_at = Some(completed);
832            }
833            if let Some(debug_msgs) = update.debug_messages {
834                step.debug_messages = Some(debug_msgs);
835            }
836            // Clearing wins over setting: an update that resolves the gate and
837            // reschedules at once would otherwise leave a live timer behind.
838            if update.clear_approval_deadline {
839                step.approval_deadline_at = None;
840            } else if let Some(deadline) = update.approval_deadline_at {
841                step.approval_deadline_at = Some(deadline);
842            }
843            if let Some(stage) = update.approval_stage {
844                step.approval_stage = stage;
845            }
846            if let Some(assignee) = update.approval_assignee {
847                step.approval_assignee = Some(assignee);
848            }
849            if let Some(requirement) = update.approval_requirement {
850                step.approval_requirement = Some(requirement);
851            }
852
853            step.updated_at = now;
854            Ok(())
855        })
856    }
857
858    fn get_step(&self, id: Uuid) -> StoreFuture<'_, Option<Step>> {
859        Box::pin(async move {
860            let state = self.state.read().await;
861            Ok(state.steps.get(&id).cloned())
862        })
863    }
864
865    fn record_step_approval(&self, step_id: Uuid, approval: StepApproval) -> StoreFuture<'_, Step> {
866        Box::pin(async move {
867            let mut state = self.state.write().await;
868            let step = state
869                .steps
870                .get_mut(&step_id)
871                .ok_or(StoreError::StepNotFound(step_id))?;
872
873            if !step.approvals.iter().any(|a| a.user_id == approval.user_id) {
874                step.approvals.push(approval);
875                step.updated_at = Utc::now();
876            }
877            Ok(step.clone())
878        })
879    }
880
881    fn list_steps(&self, run_id: Uuid) -> StoreFuture<'_, Vec<Step>> {
882        Box::pin(async move {
883            let state = self.state.read().await;
884            let mut steps: Vec<Step> = state
885                .steps
886                .values()
887                .filter(|s| s.run_id == run_id)
888                .cloned()
889                .collect();
890            steps.sort_by_key(|s| s.position);
891            Ok(steps)
892        })
893    }
894
895    fn get_stats(&self, filter: RunFilter) -> StoreFuture<'_, RunStats> {
896        Box::pin(async move {
897            let state = self.state.read().await;
898
899            let mut total_cost_usd = Decimal::ZERO;
900            let mut total_duration_ms = 0u64;
901            let mut total_runs = 0u64;
902            let mut completed_runs = 0u64;
903            let mut failed_runs = 0u64;
904            let mut cancelled_runs = 0u64;
905            let mut active_runs = 0u64;
906            let mut awaiting_approval_runs = 0u64;
907
908            for run in state.runs.values() {
909                if !run_matches_filter(run, &filter, &state.steps) {
910                    continue;
911                }
912
913                total_cost_usd += run.cost_usd;
914                total_duration_ms += run.duration_ms;
915                total_runs += 1;
916
917                match run.status.state {
918                    RunStatus::Completed | RunStatus::Warning => completed_runs += 1,
919                    RunStatus::Failed => failed_runs += 1,
920                    RunStatus::Cancelled => cancelled_runs += 1,
921                    RunStatus::AwaitingApproval => {
922                        active_runs += 1;
923                        awaiting_approval_runs += 1;
924                    }
925                    RunStatus::Pending
926                    | RunStatus::Running
927                    | RunStatus::Retrying
928                    | RunStatus::Sleeping => {
929                        active_runs += 1;
930                    }
931                }
932            }
933
934            Ok(RunStats {
935                total_runs,
936                completed_runs,
937                failed_runs,
938                cancelled_runs,
939                active_runs,
940                awaiting_approval_runs,
941                total_cost_usd,
942                total_duration_ms,
943            })
944        })
945    }
946
947    fn get_stats_history(
948        &self,
949        filter: StatsHistoryFilter,
950    ) -> StoreFuture<'_, Vec<StatsHistoryBucket>> {
951        Box::pin(async move {
952            let state = self.state.read().await;
953            let now = Utc::now();
954            let start = now - Duration::hours(filter.period.hours());
955            let run_filter = filter.to_run_filter();
956            let buckets = aggregate_history_buckets(
957                state
958                    .runs
959                    .values()
960                    .filter(|r| run_matches_filter(r, &run_filter, &state.steps)),
961                start,
962                now,
963                filter.granularity,
964            );
965            Ok(buckets)
966        })
967    }
968
969    fn create_step_dependencies(&self, deps: Vec<NewStepDependency>) -> StoreFuture<'_, ()> {
970        Box::pin(async move {
971            let mut state = self.state.write().await;
972
973            for dep in deps {
974                if !state.steps.contains_key(&dep.step_id) {
975                    return Err(StoreError::StepNotFound(dep.step_id));
976                }
977                if !state.steps.contains_key(&dep.depends_on) {
978                    return Err(StoreError::StepNotFound(dep.depends_on));
979                }
980
981                let already_exists = state
982                    .step_dependencies
983                    .iter()
984                    .any(|d| d.step_id == dep.step_id && d.depends_on == dep.depends_on);
985
986                if !already_exists {
987                    state.step_dependencies.push(StepDependency {
988                        step_id: dep.step_id,
989                        depends_on: dep.depends_on,
990                        created_at: Utc::now(),
991                    });
992                }
993            }
994
995            Ok(())
996        })
997    }
998
999    fn list_step_dependencies(&self, run_id: Uuid) -> StoreFuture<'_, Vec<StepDependency>> {
1000        Box::pin(async move {
1001            let state = self.state.read().await;
1002
1003            let run_step_ids: std::collections::HashSet<Uuid> = state
1004                .steps
1005                .values()
1006                .filter(|s| s.run_id == run_id)
1007                .map(|s| s.id)
1008                .collect();
1009
1010            let mut deps: Vec<StepDependency> = state
1011                .step_dependencies
1012                .iter()
1013                .filter(|d| run_step_ids.contains(&d.step_id))
1014                .cloned()
1015                .collect();
1016
1017            deps.sort_by_key(|d| d.created_at);
1018            Ok(deps)
1019        })
1020    }
1021}
1022
1023#[cfg(test)]
1024mod tests {
1025    use std::collections::HashMap;
1026    use std::time::Duration;
1027
1028    use chrono::TimeDelta;
1029    use serde_json::json;
1030    use tokio::spawn;
1031    use tokio::time::sleep;
1032
1033    use super::*;
1034    use crate::api_key_store::ApiKeyStore;
1035    use crate::entities::{
1036        ApiKeyScope, ApiKeyUpdate, ApprovalRequirement, NewApiKey, NewUser, TriggerKind,
1037    };
1038    use crate::user_store::UserStore;
1039
1040    use crate::memory::tests::{create_terminal_run, new_run_req};
1041    use crate::store::RunStore;
1042
1043    use crate::entities::{StepKind, step_trace_id};
1044
1045    fn new_step_req(run_id: Uuid, name: &str, position: u32) -> NewStep {
1046        NewStep {
1047            run_id,
1048            trace_id: step_trace_id(run_id, name, position),
1049            name: name.to_string(),
1050            kind: StepKind::Shell,
1051            position,
1052            input: None,
1053            is_error_handler: false,
1054        }
1055    }
1056
1057    // ---- create_run ----
1058
1059    #[tokio::test]
1060    async fn create_run_returns_pending_status() {
1061        let store = InMemoryStore::new();
1062        let run = store
1063            .create_run(new_run_req("test"))
1064            .await
1065            .unwrap()
1066            .into_run();
1067        assert_eq!(run.status.state, RunStatus::Pending);
1068        assert_eq!(run.workflow_name, "test");
1069        assert_eq!(run.retry_count, 0);
1070        assert_eq!(run.max_retries, 3);
1071    }
1072
1073    #[tokio::test]
1074    async fn create_run_generates_unique_ids() {
1075        let store = InMemoryStore::new();
1076        let r1 = store.create_run(new_run_req("a")).await.unwrap().into_run();
1077        let r2 = store.create_run(new_run_req("b")).await.unwrap().into_run();
1078        assert_ne!(r1.id, r2.id);
1079    }
1080
1081    // ---- get_run ----
1082
1083    #[tokio::test]
1084    async fn get_run_returns_created_run() {
1085        let store = InMemoryStore::new();
1086        let run = store
1087            .create_run(new_run_req("test"))
1088            .await
1089            .unwrap()
1090            .into_run();
1091        let fetched = store.get_run(run.id).await.unwrap();
1092        assert!(fetched.is_some());
1093        assert_eq!(fetched.unwrap().id, run.id);
1094    }
1095
1096    #[tokio::test]
1097    async fn get_run_returns_none_for_missing() {
1098        let store = InMemoryStore::new();
1099        let fetched = store.get_run(Uuid::nil()).await.unwrap();
1100        assert!(fetched.is_none());
1101    }
1102
1103    // ---- update_run_status ----
1104
1105    #[tokio::test]
1106    async fn update_run_status_valid_transition() {
1107        let store = InMemoryStore::new();
1108        let run = store
1109            .create_run(new_run_req("test"))
1110            .await
1111            .unwrap()
1112            .into_run();
1113
1114        store
1115            .update_run_status(run.id, RunStatus::Running)
1116            .await
1117            .unwrap();
1118
1119        let fetched = store.get_run(run.id).await.unwrap().unwrap();
1120        assert_eq!(fetched.status.state, RunStatus::Running);
1121        assert!(fetched.started_at.is_some());
1122    }
1123
1124    #[tokio::test]
1125    async fn update_run_status_invalid_transition_returns_error() {
1126        let store = InMemoryStore::new();
1127        let run = store
1128            .create_run(new_run_req("test"))
1129            .await
1130            .unwrap()
1131            .into_run();
1132
1133        let result = store.update_run_status(run.id, RunStatus::Completed).await;
1134        assert!(result.is_err());
1135
1136        let err = result.unwrap_err();
1137        assert!(matches!(err, StoreError::InvalidTransition { .. }));
1138    }
1139
1140    #[tokio::test]
1141    async fn update_run_status_not_found() {
1142        let store = InMemoryStore::new();
1143        let result = store
1144            .update_run_status(Uuid::nil(), RunStatus::Running)
1145            .await;
1146        assert!(matches!(result.unwrap_err(), StoreError::RunNotFound(_)));
1147    }
1148
1149    #[tokio::test]
1150    async fn update_run_status_terminal_sets_completed_at() {
1151        let store = InMemoryStore::new();
1152        let run = store
1153            .create_run(new_run_req("test"))
1154            .await
1155            .unwrap()
1156            .into_run();
1157
1158        store
1159            .update_run_status(run.id, RunStatus::Running)
1160            .await
1161            .unwrap();
1162        store
1163            .update_run_status(run.id, RunStatus::Completed)
1164            .await
1165            .unwrap();
1166
1167        let fetched = store.get_run(run.id).await.unwrap().unwrap();
1168        assert_eq!(fetched.status.state, RunStatus::Completed);
1169        assert!(fetched.completed_at.is_some());
1170    }
1171
1172    #[tokio::test]
1173    async fn update_run_status_terminal_to_same_is_idempotent() {
1174        let store = InMemoryStore::new();
1175        let run = store
1176            .create_run(new_run_req("test"))
1177            .await
1178            .unwrap()
1179            .into_run();
1180
1181        store
1182            .update_run_status(run.id, RunStatus::Running)
1183            .await
1184            .unwrap();
1185        store
1186            .update_run_status(run.id, RunStatus::Failed)
1187            .await
1188            .unwrap();
1189
1190        let before = store.get_run(run.id).await.unwrap().unwrap();
1191        let completed_at_before = before.completed_at;
1192
1193        store
1194            .update_run_status(run.id, RunStatus::Failed)
1195            .await
1196            .unwrap();
1197
1198        let after = store.get_run(run.id).await.unwrap().unwrap();
1199        assert_eq!(after.status.state, RunStatus::Failed);
1200        assert_eq!(after.completed_at, completed_at_before);
1201    }
1202
1203    #[tokio::test]
1204    async fn update_run_terminal_to_same_via_update_run_is_idempotent() {
1205        let store = InMemoryStore::new();
1206        let run = store
1207            .create_run(new_run_req("test"))
1208            .await
1209            .unwrap()
1210            .into_run();
1211
1212        store
1213            .update_run_status(run.id, RunStatus::Running)
1214            .await
1215            .unwrap();
1216        store
1217            .update_run(
1218                run.id,
1219                RunUpdate {
1220                    status: Some(RunStatus::Failed),
1221                    error: Some("first failure".to_string()),
1222                    ..RunUpdate::default()
1223                },
1224            )
1225            .await
1226            .unwrap();
1227
1228        let before = store.get_run(run.id).await.unwrap().unwrap();
1229
1230        store
1231            .update_run(
1232                run.id,
1233                RunUpdate {
1234                    status: Some(RunStatus::Failed),
1235                    ..RunUpdate::default()
1236                },
1237            )
1238            .await
1239            .unwrap();
1240
1241        let after = store.get_run(run.id).await.unwrap().unwrap();
1242        assert_eq!(after.status.state, RunStatus::Failed);
1243        assert_eq!(after.completed_at, before.completed_at);
1244        assert_eq!(after.error, Some("first failure".to_string()));
1245    }
1246
1247    // ---- list_runs ----
1248
1249    #[tokio::test]
1250    async fn list_runs_empty_store() {
1251        let store = InMemoryStore::new();
1252        let page = store.list_runs(RunFilter::default(), 1, 20).await.unwrap();
1253        assert_eq!(page.total, 0);
1254        assert!(page.items.is_empty());
1255    }
1256
1257    #[tokio::test]
1258    async fn list_runs_with_workflow_filter() {
1259        let store = InMemoryStore::new();
1260        store
1261            .create_run(new_run_req("deploy"))
1262            .await
1263            .unwrap()
1264            .into_run();
1265        store
1266            .create_run(new_run_req("test"))
1267            .await
1268            .unwrap()
1269            .into_run();
1270        store
1271            .create_run(new_run_req("deploy"))
1272            .await
1273            .unwrap()
1274            .into_run();
1275
1276        let filter = RunFilter {
1277            workflow_name: Some("deploy".to_string()),
1278            ..RunFilter::default()
1279        };
1280        let page = store.list_runs(filter, 1, 20).await.unwrap();
1281        assert_eq!(page.total, 2);
1282        assert!(page.items.iter().all(|r| r.workflow_name == "deploy"));
1283    }
1284
1285    #[tokio::test]
1286    async fn list_runs_with_status_filter() {
1287        let store = InMemoryStore::new();
1288        let run = store.create_run(new_run_req("a")).await.unwrap().into_run();
1289        store.create_run(new_run_req("b")).await.unwrap().into_run();
1290
1291        store
1292            .update_run_status(run.id, RunStatus::Running)
1293            .await
1294            .unwrap();
1295
1296        let filter = RunFilter {
1297            status: Some(RunStatus::Running),
1298            ..RunFilter::default()
1299        };
1300        let page = store.list_runs(filter, 1, 20).await.unwrap();
1301        assert_eq!(page.total, 1);
1302        assert_eq!(page.items[0].id, run.id);
1303    }
1304
1305    #[tokio::test]
1306    async fn list_runs_pagination() {
1307        let store = InMemoryStore::new();
1308        for i in 0..5 {
1309            store
1310                .create_run(new_run_req(&format!("wf-{i}")))
1311                .await
1312                .unwrap()
1313                .into_run();
1314        }
1315
1316        let page1 = store.list_runs(RunFilter::default(), 1, 2).await.unwrap();
1317        assert_eq!(page1.total, 5);
1318        assert_eq!(page1.items.len(), 2);
1319        assert_eq!(page1.page, 1);
1320        assert_eq!(page1.per_page, 2);
1321
1322        let page2 = store.list_runs(RunFilter::default(), 2, 2).await.unwrap();
1323        assert_eq!(page2.items.len(), 2);
1324
1325        let page3 = store.list_runs(RunFilter::default(), 3, 2).await.unwrap();
1326        assert_eq!(page3.items.len(), 1);
1327    }
1328
1329    // ---- worker lease ----
1330
1331    fn lease(worker_id: &str, ttl_secs: u64) -> Option<LeaseRequest> {
1332        Some(LeaseRequest {
1333            worker_id: worker_id.to_string(),
1334            ttl: Duration::from_secs(ttl_secs),
1335        })
1336    }
1337
1338    /// Pick a run holding a lease that is already in the past.
1339    ///
1340    /// The short sleep matters: `Utc::now()` has microsecond resolution, so a
1341    /// sub-microsecond TTL can still read as "not yet expired" in the same tick.
1342    async fn pick_with_expired_lease(store: &InMemoryStore, max_retries: u32) -> Run {
1343        let mut req = new_run_req("test");
1344        req.max_retries = max_retries;
1345        store.create_run(req).await.unwrap();
1346        let picked = expire_now(store).await;
1347        picked.expect("a pending run was just created")
1348    }
1349
1350    /// Pick the next pending run with a lease that is expired by the time this
1351    /// returns.
1352    async fn expire_now(store: &InMemoryStore) -> Option<Run> {
1353        let picked = store
1354            .pick_next_pending(Some(LeaseRequest {
1355                worker_id: "worker-1".to_string(),
1356                ttl: Duration::from_nanos(1),
1357            }))
1358            .await
1359            .unwrap();
1360        sleep(Duration::from_millis(2)).await;
1361        picked
1362    }
1363
1364    #[tokio::test]
1365    async fn pick_next_pending_attaches_lease() {
1366        let store = InMemoryStore::new();
1367        store.create_run(new_run_req("test")).await.unwrap();
1368
1369        let picked = store
1370            .pick_next_pending(lease("worker-1", 90))
1371            .await
1372            .unwrap()
1373            .unwrap();
1374
1375        assert_eq!(picked.worker_id.as_deref(), Some("worker-1"));
1376        let expires = picked.lease_expires_at.expect("lease set");
1377        assert!(expires > Utc::now());
1378        assert!(expires <= Utc::now() + TimeDelta::seconds(91));
1379    }
1380
1381    #[tokio::test]
1382    async fn pick_next_pending_without_lease_leaves_run_unowned() {
1383        let store = InMemoryStore::new();
1384        store.create_run(new_run_req("test")).await.unwrap();
1385
1386        let picked = store.pick_next_pending(None).await.unwrap().unwrap();
1387
1388        assert!(picked.worker_id.is_none());
1389        assert!(picked.lease_expires_at.is_none());
1390    }
1391
1392    #[tokio::test]
1393    async fn renew_lease_extends_expiry_for_owner() {
1394        let store = InMemoryStore::new();
1395        store.create_run(new_run_req("test")).await.unwrap();
1396        let picked = store
1397            .pick_next_pending(lease("worker-1", 1))
1398            .await
1399            .unwrap()
1400            .unwrap();
1401
1402        let renewed = store
1403            .renew_lease(picked.id, lease("worker-1", 90).unwrap())
1404            .await
1405            .unwrap();
1406
1407        assert!(renewed > picked.lease_expires_at.unwrap());
1408        let after = store.get_run(picked.id).await.unwrap().unwrap();
1409        assert_eq!(after.lease_expires_at, Some(renewed));
1410    }
1411
1412    #[tokio::test]
1413    async fn renew_lease_rejects_other_worker() {
1414        let store = InMemoryStore::new();
1415        store.create_run(new_run_req("test")).await.unwrap();
1416        let picked = store
1417            .pick_next_pending(lease("worker-1", 90))
1418            .await
1419            .unwrap()
1420            .unwrap();
1421
1422        let err = store
1423            .renew_lease(picked.id, lease("worker-2", 90).unwrap())
1424            .await
1425            .unwrap_err();
1426
1427        assert!(matches!(
1428            err,
1429            StoreError::LeaseLost { held_by: Some(ref w), .. } if w == "worker-1"
1430        ));
1431    }
1432
1433    #[tokio::test]
1434    async fn renew_lease_rejects_run_that_left_running() {
1435        let store = InMemoryStore::new();
1436        store.create_run(new_run_req("test")).await.unwrap();
1437        let picked = store
1438            .pick_next_pending(lease("worker-1", 90))
1439            .await
1440            .unwrap()
1441            .unwrap();
1442        store
1443            .update_run_status(picked.id, RunStatus::Cancelled)
1444            .await
1445            .unwrap();
1446
1447        let err = store
1448            .renew_lease(picked.id, lease("worker-1", 90).unwrap())
1449            .await
1450            .unwrap_err();
1451
1452        assert!(matches!(err, StoreError::LeaseLost { .. }));
1453    }
1454
1455    #[tokio::test]
1456    async fn renew_lease_on_unknown_run_is_not_found() {
1457        let store = InMemoryStore::new();
1458
1459        let err = store
1460            .renew_lease(Uuid::now_v7(), lease("worker-1", 90).unwrap())
1461            .await
1462            .unwrap_err();
1463
1464        assert!(matches!(err, StoreError::RunNotFound(_)));
1465    }
1466
1467    #[tokio::test]
1468    async fn leaving_running_clears_the_lease() {
1469        for target in [
1470            RunStatus::Completed,
1471            RunStatus::Retrying,
1472            RunStatus::AwaitingApproval,
1473        ] {
1474            let store = InMemoryStore::new();
1475            store.create_run(new_run_req("test")).await.unwrap();
1476            let picked = store
1477                .pick_next_pending(lease("worker-1", 90))
1478                .await
1479                .unwrap()
1480                .unwrap();
1481
1482            store.update_run_status(picked.id, target).await.unwrap();
1483
1484            let after = store.get_run(picked.id).await.unwrap().unwrap();
1485            assert!(after.worker_id.is_none(), "worker_id kept for {target}");
1486            assert!(
1487                after.lease_expires_at.is_none(),
1488                "lease_expires_at kept for {target}"
1489            );
1490        }
1491    }
1492
1493    #[tokio::test]
1494    async fn update_run_to_terminal_clears_the_lease() {
1495        let store = InMemoryStore::new();
1496        store.create_run(new_run_req("test")).await.unwrap();
1497        let picked = store
1498            .pick_next_pending(lease("worker-1", 90))
1499            .await
1500            .unwrap()
1501            .unwrap();
1502
1503        store
1504            .update_run(
1505                picked.id,
1506                RunUpdate {
1507                    status: Some(RunStatus::Failed),
1508                    ..RunUpdate::default()
1509                },
1510            )
1511            .await
1512            .unwrap();
1513
1514        let after = store.get_run(picked.id).await.unwrap().unwrap();
1515        assert!(after.worker_id.is_none());
1516        assert!(after.lease_expires_at.is_none());
1517    }
1518
1519    #[tokio::test]
1520    async fn update_run_lease_set_with_running_status_attaches_lease() {
1521        let store = InMemoryStore::new();
1522        let run = store
1523            .create_run(new_run_req("test"))
1524            .await
1525            .unwrap()
1526            .into_run();
1527        let expires_at = Utc::now() + TimeDelta::seconds(60);
1528
1529        store
1530            .update_run(
1531                run.id,
1532                RunUpdate {
1533                    status: Some(RunStatus::Running),
1534                    lease: Some(LeaseUpdate::Set {
1535                        worker_id: "worker-1".to_string(),
1536                        expires_at,
1537                    }),
1538                    ..RunUpdate::default()
1539                },
1540            )
1541            .await
1542            .unwrap();
1543
1544        let after = store.get_run(run.id).await.unwrap().unwrap();
1545        assert_eq!(after.status.state, RunStatus::Running);
1546        assert_eq!(after.worker_id.as_deref(), Some("worker-1"));
1547        assert_eq!(after.lease_expires_at, Some(expires_at));
1548        // The transferred lease is renewable by its holder like a picked one.
1549        store
1550            .renew_lease(run.id, lease("worker-1", 90).unwrap())
1551            .await
1552            .unwrap();
1553    }
1554
1555    #[tokio::test]
1556    async fn update_run_lease_release_keeps_run_running_without_lease() {
1557        let store = InMemoryStore::new();
1558        store.create_run(new_run_req("test")).await.unwrap();
1559        let picked = store
1560            .pick_next_pending(lease("worker-1", 90))
1561            .await
1562            .unwrap()
1563            .unwrap();
1564
1565        store
1566            .update_run(
1567                picked.id,
1568                RunUpdate {
1569                    lease: Some(LeaseUpdate::Release),
1570                    ..RunUpdate::default()
1571                },
1572            )
1573            .await
1574            .unwrap();
1575
1576        let after = store.get_run(picked.id).await.unwrap().unwrap();
1577        assert_eq!(after.status.state, RunStatus::Running);
1578        assert!(after.worker_id.is_none());
1579        assert!(after.lease_expires_at.is_none());
1580        let err = store
1581            .renew_lease(picked.id, lease("worker-1", 90).unwrap())
1582            .await
1583            .unwrap_err();
1584        assert!(matches!(err, StoreError::LeaseLost { .. }));
1585    }
1586
1587    #[tokio::test]
1588    async fn update_run_lease_none_leaves_lease_untouched() {
1589        let store = InMemoryStore::new();
1590        store.create_run(new_run_req("test")).await.unwrap();
1591        let picked = store
1592            .pick_next_pending(lease("worker-1", 90))
1593            .await
1594            .unwrap()
1595            .unwrap();
1596
1597        store
1598            .update_run(
1599                picked.id,
1600                RunUpdate {
1601                    cost_usd: Some(Decimal::new(150, 2)),
1602                    ..RunUpdate::default()
1603                },
1604            )
1605            .await
1606            .unwrap();
1607
1608        let after = store.get_run(picked.id).await.unwrap().unwrap();
1609        assert_eq!(after.worker_id.as_deref(), Some("worker-1"));
1610        assert_eq!(after.lease_expires_at, picked.lease_expires_at);
1611    }
1612
1613    // ---- reap_expired_leases ----
1614
1615    #[tokio::test]
1616    async fn reap_expired_leases_empty_store() {
1617        let store = InMemoryStore::new();
1618        assert!(store.reap_expired_leases(100).await.unwrap().is_empty());
1619    }
1620
1621    #[tokio::test]
1622    async fn reap_expired_leases_requeues_run() {
1623        let store = InMemoryStore::new();
1624        let picked = pick_with_expired_lease(&store, 3).await;
1625
1626        let reaped = store.reap_expired_leases(100).await.unwrap();
1627
1628        assert_eq!(reaped.len(), 1);
1629        assert_eq!(reaped[0].from, RunStatus::Running);
1630        assert_eq!(reaped[0].to, RunStatus::Pending);
1631
1632        let after = store.get_run(picked.id).await.unwrap().unwrap();
1633        assert_eq!(after.status.state, RunStatus::Pending);
1634        assert_eq!(after.retry_count, 0);
1635        assert_eq!(after.lease_recoveries, 1);
1636        assert!(after.error.is_none());
1637        assert!(after.worker_id.is_none());
1638        assert!(after.lease_expires_at.is_none());
1639    }
1640
1641    #[tokio::test]
1642    async fn reap_expired_leases_requeued_run_is_pickable_again() {
1643        let store = InMemoryStore::new();
1644        let picked = pick_with_expired_lease(&store, 3).await;
1645        store.reap_expired_leases(100).await.unwrap();
1646
1647        let repicked = store
1648            .pick_next_pending(lease("worker-2", 90))
1649            .await
1650            .unwrap()
1651            .unwrap();
1652
1653        assert_eq!(repicked.id, picked.id);
1654        assert_eq!(repicked.worker_id.as_deref(), Some("worker-2"));
1655    }
1656
1657    #[tokio::test]
1658    async fn reap_expired_leases_ignores_valid_lease() {
1659        let store = InMemoryStore::new();
1660        store.create_run(new_run_req("test")).await.unwrap();
1661        let picked = store
1662            .pick_next_pending(lease("worker-1", 90))
1663            .await
1664            .unwrap()
1665            .unwrap();
1666
1667        assert!(store.reap_expired_leases(100).await.unwrap().is_empty());
1668
1669        let after = store.get_run(picked.id).await.unwrap().unwrap();
1670        assert_eq!(after.status.state, RunStatus::Running);
1671        assert_eq!(after.retry_count, 0);
1672    }
1673
1674    #[tokio::test]
1675    async fn reap_expired_leases_ignores_run_without_lease() {
1676        let store = InMemoryStore::new();
1677        store.create_run(new_run_req("test")).await.unwrap();
1678        let picked = store.pick_next_pending(None).await.unwrap().unwrap();
1679
1680        assert!(store.reap_expired_leases(100).await.unwrap().is_empty());
1681
1682        let after = store.get_run(picked.id).await.unwrap().unwrap();
1683        assert_eq!(after.status.state, RunStatus::Running);
1684    }
1685
1686    #[tokio::test]
1687    async fn reap_expired_leases_fails_run_when_retries_exhausted() {
1688        let store = InMemoryStore::new();
1689        let picked = pick_with_expired_lease(&store, 0).await;
1690
1691        let reaped = store.reap_expired_leases(100).await.unwrap();
1692
1693        assert_eq!(reaped[0].to, RunStatus::Failed);
1694        let after = store.get_run(picked.id).await.unwrap().unwrap();
1695        assert_eq!(after.status.state, RunStatus::Failed);
1696        assert_eq!(after.error.as_deref(), Some(LEASE_EXPIRED_ERROR));
1697        assert!(after.completed_at.is_some());
1698    }
1699
1700    #[tokio::test]
1701    async fn reap_expired_leases_fails_after_max_retries_recoveries() {
1702        let store = InMemoryStore::new();
1703        let picked = pick_with_expired_lease(&store, 2).await;
1704
1705        // Recovery 1 and 2 requeue, recovery 3 gives up.
1706        for expected in [RunStatus::Pending, RunStatus::Pending, RunStatus::Failed] {
1707            let reaped = store.reap_expired_leases(100).await.unwrap();
1708            assert_eq!(reaped[0].to, expected);
1709            if expected == RunStatus::Pending {
1710                expire_now(&store).await;
1711            }
1712        }
1713
1714        let after = store.get_run(picked.id).await.unwrap().unwrap();
1715        assert_eq!(after.lease_recoveries, 3);
1716        assert_eq!(after.retry_count, 0);
1717        assert_eq!(after.error.as_deref(), Some(LEASE_EXPIRED_ERROR));
1718    }
1719
1720    #[tokio::test]
1721    async fn reap_expired_leases_keeps_the_attempt_number() {
1722        let store = InMemoryStore::new();
1723        let picked = pick_with_expired_lease(&store, 3).await;
1724        let before = store
1725            .create_step(new_step_req(picked.id, "build", 0))
1726            .await
1727            .unwrap();
1728
1729        store.reap_expired_leases(100).await.unwrap();
1730        let repicked = store
1731            .pick_next_pending(lease("worker-2", 90))
1732            .await
1733            .unwrap()
1734            .unwrap();
1735        let after = store
1736            .create_step(new_step_req(repicked.id, "build", 0))
1737            .await
1738            .unwrap();
1739
1740        assert_eq!(before.attempt, 1);
1741        assert_eq!(after.attempt, before.attempt);
1742        assert_eq!(repicked.retry_count, 0);
1743        assert_eq!(repicked.lease_recoveries, 1);
1744    }
1745
1746    #[tokio::test]
1747    async fn reap_expired_leases_counts_apart_from_handler_retries() {
1748        let store = InMemoryStore::new();
1749        let picked = pick_with_expired_lease(&store, 1).await;
1750        store
1751            .update_run(
1752                picked.id,
1753                RunUpdate {
1754                    increment_retry: true,
1755                    ..RunUpdate::default()
1756                },
1757            )
1758            .await
1759            .unwrap();
1760
1761        // One handler retry already consumed max_retries: the lease recovery
1762        // budget is separate, so the first reap still requeues.
1763        let reaped = store.reap_expired_leases(100).await.unwrap();
1764
1765        assert_eq!(reaped[0].to, RunStatus::Pending);
1766        assert_eq!(reaped[0].run.retry_count, 1);
1767        assert_eq!(reaped[0].run.lease_recoveries, 1);
1768    }
1769
1770    #[tokio::test]
1771    async fn reap_expired_leases_respects_limit() {
1772        let store = InMemoryStore::new();
1773        for _ in 0..3 {
1774            pick_with_expired_lease(&store, 3).await;
1775        }
1776
1777        let reaped = store.reap_expired_leases(2).await.unwrap();
1778        assert_eq!(reaped.len(), 2);
1779
1780        let rest = store.reap_expired_leases(100).await.unwrap();
1781        assert_eq!(rest.len(), 1);
1782    }
1783
1784    // ---- pick_next_pending ----
1785
1786    #[tokio::test]
1787    async fn pick_next_pending_empty_store() {
1788        let store = InMemoryStore::new();
1789        let result = store.pick_next_pending(None).await.unwrap();
1790        assert!(result.is_none());
1791    }
1792
1793    #[tokio::test]
1794    async fn pick_next_pending_returns_oldest_and_transitions_to_running() {
1795        let store = InMemoryStore::new();
1796        let r1 = store
1797            .create_run(new_run_req("first"))
1798            .await
1799            .unwrap()
1800            .into_run();
1801        let _r2 = store
1802            .create_run(new_run_req("second"))
1803            .await
1804            .unwrap()
1805            .into_run();
1806
1807        let picked = store.pick_next_pending(None).await.unwrap().unwrap();
1808        assert_eq!(picked.id, r1.id);
1809        assert_eq!(picked.status.state, RunStatus::Running);
1810        assert!(picked.started_at.is_some());
1811
1812        // Verify it's Running in the store too.
1813        let fetched = store.get_run(r1.id).await.unwrap().unwrap();
1814        assert_eq!(fetched.status.state, RunStatus::Running);
1815    }
1816
1817    #[tokio::test]
1818    async fn pick_next_pending_skips_non_pending() {
1819        let store = InMemoryStore::new();
1820        let r1 = store.create_run(new_run_req("a")).await.unwrap().into_run();
1821        let r2 = store.create_run(new_run_req("b")).await.unwrap().into_run();
1822
1823        // Transition r1 to Running.
1824        store
1825            .update_run_status(r1.id, RunStatus::Running)
1826            .await
1827            .unwrap();
1828
1829        let picked = store.pick_next_pending(None).await.unwrap().unwrap();
1830        assert_eq!(picked.id, r2.id);
1831    }
1832
1833    // ---- create_step ----
1834
1835    #[tokio::test]
1836    async fn create_step_returns_pending() {
1837        let store = InMemoryStore::new();
1838        let run = store
1839            .create_run(new_run_req("test"))
1840            .await
1841            .unwrap()
1842            .into_run();
1843
1844        let step = store
1845            .create_step(NewStep {
1846                run_id: run.id,
1847                trace_id: step_trace_id(run.id, "build", 0),
1848                name: "build".to_string(),
1849                kind: crate::entities::StepKind::Shell,
1850                position: 0,
1851                input: Some(json!({"command": "cargo build"})),
1852                is_error_handler: false,
1853            })
1854            .await
1855            .unwrap();
1856
1857        assert_eq!(step.status.state, StepStatus::Pending);
1858        assert_eq!(step.name, "build");
1859        assert_eq!(step.run_id, run.id);
1860        assert_eq!(step.position, 0);
1861    }
1862
1863    #[tokio::test]
1864    async fn create_step_for_missing_run_returns_error() {
1865        let store = InMemoryStore::new();
1866        let result = store
1867            .create_step(NewStep {
1868                run_id: Uuid::nil(),
1869                trace_id: step_trace_id(Uuid::nil(), "build", 0),
1870                name: "build".to_string(),
1871                kind: crate::entities::StepKind::Shell,
1872                position: 0,
1873                input: None,
1874                is_error_handler: false,
1875            })
1876            .await;
1877        assert!(matches!(result.unwrap_err(), StoreError::RunNotFound(_)));
1878    }
1879
1880    // ---- update_step ----
1881
1882    #[tokio::test]
1883    async fn update_step_applies_partial_update() {
1884        let store = InMemoryStore::new();
1885        let run = store
1886            .create_run(new_run_req("test"))
1887            .await
1888            .unwrap()
1889            .into_run();
1890
1891        let step = store
1892            .create_step(NewStep {
1893                run_id: run.id,
1894                trace_id: step_trace_id(run.id, "build", 0),
1895                name: "build".to_string(),
1896                kind: crate::entities::StepKind::Shell,
1897                position: 0,
1898                input: None,
1899                is_error_handler: false,
1900            })
1901            .await
1902            .unwrap();
1903
1904        // Transition Pending → Running first
1905        store
1906            .update_step(
1907                step.id,
1908                StepUpdate {
1909                    status: Some(StepStatus::Running),
1910                    ..StepUpdate::default()
1911                },
1912            )
1913            .await
1914            .unwrap();
1915
1916        // Then Running → Completed
1917        store
1918            .update_step(
1919                step.id,
1920                StepUpdate {
1921                    status: Some(StepStatus::Completed),
1922                    output: Some(json!({"stdout": "ok"})),
1923                    duration_ms: Some(150),
1924                    ..StepUpdate::default()
1925                },
1926            )
1927            .await
1928            .unwrap();
1929
1930        let steps = store.list_steps(run.id).await.unwrap();
1931        assert_eq!(steps.len(), 1);
1932        assert_eq!(steps[0].status.state, StepStatus::Completed);
1933        assert_eq!(steps[0].duration_ms, 150);
1934        assert!(steps[0].output.is_some());
1935    }
1936
1937    #[tokio::test]
1938    async fn update_step_not_found() {
1939        let store = InMemoryStore::new();
1940        let result = store.update_step(Uuid::nil(), StepUpdate::default()).await;
1941        assert!(matches!(result.unwrap_err(), StoreError::StepNotFound(_)));
1942    }
1943
1944    // ---- list_steps ----
1945
1946    #[tokio::test]
1947    async fn list_steps_ordered_by_position() {
1948        let store = InMemoryStore::new();
1949        let run = store
1950            .create_run(new_run_req("test"))
1951            .await
1952            .unwrap()
1953            .into_run();
1954
1955        // Insert out of order.
1956        store
1957            .create_step(NewStep {
1958                run_id: run.id,
1959                trace_id: step_trace_id(run.id, "deploy", 2),
1960                name: "deploy".to_string(),
1961                kind: crate::entities::StepKind::Shell,
1962                position: 2,
1963                input: None,
1964                is_error_handler: false,
1965            })
1966            .await
1967            .unwrap();
1968        store
1969            .create_step(NewStep {
1970                run_id: run.id,
1971                trace_id: step_trace_id(run.id, "build", 0),
1972                name: "build".to_string(),
1973                kind: crate::entities::StepKind::Shell,
1974                position: 0,
1975                input: None,
1976                is_error_handler: false,
1977            })
1978            .await
1979            .unwrap();
1980        store
1981            .create_step(NewStep {
1982                run_id: run.id,
1983                trace_id: step_trace_id(run.id, "test", 1),
1984                name: "test".to_string(),
1985                kind: crate::entities::StepKind::Shell,
1986                position: 1,
1987                input: None,
1988                is_error_handler: false,
1989            })
1990            .await
1991            .unwrap();
1992
1993        let steps = store.list_steps(run.id).await.unwrap();
1994        assert_eq!(steps.len(), 3);
1995        assert_eq!(steps[0].name, "build");
1996        assert_eq!(steps[1].name, "test");
1997        assert_eq!(steps[2].name, "deploy");
1998    }
1999
2000    #[tokio::test]
2001    async fn list_steps_empty_for_run_without_steps() {
2002        let store = InMemoryStore::new();
2003        let run = store
2004            .create_run(new_run_req("test"))
2005            .await
2006            .unwrap()
2007            .into_run();
2008        let steps = store.list_steps(run.id).await.unwrap();
2009        assert!(steps.is_empty());
2010    }
2011
2012    // ---- update_run ----
2013
2014    #[tokio::test]
2015    async fn update_run_applies_cost_and_duration() {
2016        let store = InMemoryStore::new();
2017        let run = store
2018            .create_run(new_run_req("test"))
2019            .await
2020            .unwrap()
2021            .into_run();
2022
2023        store
2024            .update_run(
2025                run.id,
2026                RunUpdate {
2027                    cost_usd: Some(Decimal::new(123, 2)),
2028                    duration_ms: Some(5000),
2029                    ..RunUpdate::default()
2030                },
2031            )
2032            .await
2033            .unwrap();
2034
2035        let fetched = store.get_run(run.id).await.unwrap().unwrap();
2036        assert_eq!(fetched.cost_usd, Decimal::new(123, 2));
2037        assert_eq!(fetched.duration_ms, 5000);
2038    }
2039
2040    #[tokio::test]
2041    async fn update_run_increment_retry() {
2042        let store = InMemoryStore::new();
2043        let run = store
2044            .create_run(new_run_req("test"))
2045            .await
2046            .unwrap()
2047            .into_run();
2048        assert_eq!(run.retry_count, 0);
2049
2050        store
2051            .update_run(
2052                run.id,
2053                RunUpdate {
2054                    increment_retry: true,
2055                    ..RunUpdate::default()
2056                },
2057            )
2058            .await
2059            .unwrap();
2060
2061        let fetched = store.get_run(run.id).await.unwrap().unwrap();
2062        assert_eq!(fetched.retry_count, 1);
2063    }
2064
2065    #[tokio::test]
2066    async fn update_run_not_found() {
2067        let store = InMemoryStore::new();
2068        let result = store.update_run(Uuid::nil(), RunUpdate::default()).await;
2069        assert!(matches!(result.unwrap_err(), StoreError::RunNotFound(_)));
2070    }
2071
2072    // ---- concurrent access ----
2073
2074    #[tokio::test]
2075    async fn concurrent_pick_next_pending_no_double_pick() {
2076        let store = InMemoryStore::new();
2077
2078        // Create 10 pending runs.
2079        for i in 0..10 {
2080            store
2081                .create_run(new_run_req(&format!("wf-{i}")))
2082                .await
2083                .unwrap()
2084                .into_run();
2085        }
2086
2087        // Concurrently pick from multiple tasks.
2088        let mut handles = Vec::new();
2089        for _ in 0..10 {
2090            let s = store.clone();
2091            handles.push(spawn(async move { s.pick_next_pending(None).await }));
2092        }
2093
2094        let mut picked_ids = Vec::new();
2095        for h in handles {
2096            if let Ok(Ok(Some(run))) = h.await {
2097                picked_ids.push(run.id);
2098            }
2099        }
2100
2101        // Each run should be picked at most once.
2102        let unique: std::collections::HashSet<_> = picked_ids.iter().collect();
2103        assert_eq!(unique.len(), picked_ids.len());
2104    }
2105
2106    // ---- get_stats ----
2107
2108    #[tokio::test]
2109    async fn get_stats_empty_store() {
2110        let store = InMemoryStore::new();
2111        let stats = store.get_stats(RunFilter::default()).await.unwrap();
2112        assert_eq!(stats.total_runs, 0);
2113        assert_eq!(stats.completed_runs, 0);
2114        assert_eq!(stats.failed_runs, 0);
2115        assert_eq!(stats.cancelled_runs, 0);
2116        assert_eq!(stats.active_runs, 0);
2117        assert_eq!(stats.total_cost_usd, Decimal::ZERO);
2118        assert_eq!(stats.total_duration_ms, 0);
2119    }
2120
2121    #[tokio::test]
2122    async fn get_stats_aggregates_counts_and_totals() {
2123        let store = InMemoryStore::new();
2124
2125        // Create runs in various states.
2126        let r1 = store
2127            .create_run(new_run_req("wf1"))
2128            .await
2129            .unwrap()
2130            .into_run();
2131        let r2 = store
2132            .create_run(new_run_req("wf2"))
2133            .await
2134            .unwrap()
2135            .into_run();
2136        let r3 = store
2137            .create_run(new_run_req("wf3"))
2138            .await
2139            .unwrap()
2140            .into_run();
2141        let _r4 = store
2142            .create_run(new_run_req("wf4"))
2143            .await
2144            .unwrap()
2145            .into_run();
2146
2147        // Transition r1 to Running, then Completed.
2148        store
2149            .update_run_status(r1.id, RunStatus::Running)
2150            .await
2151            .unwrap();
2152        store
2153            .update_run_status(r1.id, RunStatus::Completed)
2154            .await
2155            .unwrap();
2156
2157        // Transition r2 to Running, then Failed.
2158        store
2159            .update_run_status(r2.id, RunStatus::Running)
2160            .await
2161            .unwrap();
2162        store
2163            .update_run_status(r2.id, RunStatus::Failed)
2164            .await
2165            .unwrap();
2166
2167        // Transition r3 to Cancelled.
2168        store
2169            .update_run_status(r3.id, RunStatus::Cancelled)
2170            .await
2171            .unwrap();
2172
2173        // r4 remains Pending (active).
2174
2175        // Update cost and duration.
2176        store
2177            .update_run(
2178                r1.id,
2179                RunUpdate {
2180                    cost_usd: Some(Decimal::new(1000, 2)),
2181                    duration_ms: Some(1000),
2182                    ..RunUpdate::default()
2183                },
2184            )
2185            .await
2186            .unwrap();
2187
2188        store
2189            .update_run(
2190                r2.id,
2191                RunUpdate {
2192                    cost_usd: Some(Decimal::new(500, 2)),
2193                    duration_ms: Some(500),
2194                    ..RunUpdate::default()
2195                },
2196            )
2197            .await
2198            .unwrap();
2199
2200        let stats = store.get_stats(RunFilter::default()).await.unwrap();
2201        assert_eq!(stats.total_runs, 4);
2202        assert_eq!(stats.completed_runs, 1);
2203        assert_eq!(stats.failed_runs, 1);
2204        assert_eq!(stats.cancelled_runs, 1);
2205        assert_eq!(stats.active_runs, 1); // r4 is Pending
2206        assert_eq!(stats.total_cost_usd, Decimal::new(1500, 2));
2207        assert_eq!(stats.total_duration_ms, 1500);
2208    }
2209
2210    #[tokio::test]
2211    async fn update_run_status_running_to_retrying() {
2212        let store = InMemoryStore::new();
2213        let run = store
2214            .create_run(new_run_req("test"))
2215            .await
2216            .unwrap()
2217            .into_run();
2218
2219        store
2220            .update_run_status(run.id, RunStatus::Running)
2221            .await
2222            .unwrap();
2223
2224        store
2225            .update_run_status(run.id, RunStatus::Retrying)
2226            .await
2227            .unwrap();
2228
2229        let fetched = store.get_run(run.id).await.unwrap().unwrap();
2230        assert_eq!(fetched.status.state, RunStatus::Retrying);
2231        assert!(!fetched.status.state.is_terminal());
2232        assert!(fetched.completed_at.is_none()); // Not a terminal state
2233    }
2234
2235    #[tokio::test]
2236    async fn update_run_status_retrying_to_running_allowed() {
2237        let store = InMemoryStore::new();
2238        let run = store
2239            .create_run(new_run_req("test"))
2240            .await
2241            .unwrap()
2242            .into_run();
2243
2244        store
2245            .update_run_status(run.id, RunStatus::Running)
2246            .await
2247            .unwrap();
2248        store
2249            .update_run_status(run.id, RunStatus::Retrying)
2250            .await
2251            .unwrap();
2252
2253        // Retrying → Running should be allowed
2254        store
2255            .update_run_status(run.id, RunStatus::Running)
2256            .await
2257            .unwrap();
2258
2259        let fetched = store.get_run(run.id).await.unwrap().unwrap();
2260        assert_eq!(fetched.status.state, RunStatus::Running);
2261    }
2262
2263    #[tokio::test]
2264    async fn update_run_with_invalid_status_transition_errors() {
2265        let store = InMemoryStore::new();
2266        let run = store
2267            .create_run(new_run_req("test"))
2268            .await
2269            .unwrap()
2270            .into_run();
2271
2272        // Try to apply invalid status transition via update_run
2273        let result = store
2274            .update_run(
2275                run.id,
2276                RunUpdate {
2277                    status: Some(RunStatus::Completed), // invalid from Pending
2278                    ..RunUpdate::default()
2279                },
2280            )
2281            .await;
2282
2283        assert!(result.is_err());
2284    }
2285
2286    #[tokio::test]
2287    async fn create_step_with_complex_input() {
2288        let store = InMemoryStore::new();
2289        let run = store
2290            .create_run(new_run_req("test"))
2291            .await
2292            .unwrap()
2293            .into_run();
2294
2295        let complex_input = json!({
2296            "command": "cargo build",
2297            "env": {
2298                "RUST_LOG": "debug",
2299                "CUSTOM": "value"
2300            },
2301            "timeout": 60,
2302            "retry_policy": {
2303                "max_attempts": 3,
2304                "backoff": "exponential"
2305            }
2306        });
2307
2308        let step = store
2309            .create_step(NewStep {
2310                run_id: run.id,
2311                trace_id: step_trace_id(run.id, "build", 0),
2312                name: "build".to_string(),
2313                kind: crate::entities::StepKind::Agent,
2314                position: 0,
2315                input: Some(complex_input.clone()),
2316                is_error_handler: false,
2317            })
2318            .await
2319            .unwrap();
2320
2321        assert_eq!(step.input, Some(complex_input));
2322    }
2323
2324    #[tokio::test]
2325    async fn update_step_with_error_message() {
2326        let store = InMemoryStore::new();
2327        let run = store
2328            .create_run(new_run_req("test"))
2329            .await
2330            .unwrap()
2331            .into_run();
2332
2333        let step = store
2334            .create_step(NewStep {
2335                run_id: run.id,
2336                trace_id: step_trace_id(run.id, "build", 0),
2337                name: "build".to_string(),
2338                kind: crate::entities::StepKind::Shell,
2339                position: 0,
2340                input: None,
2341                is_error_handler: false,
2342            })
2343            .await
2344            .unwrap();
2345
2346        store
2347            .update_step(
2348                step.id,
2349                StepUpdate {
2350                    status: Some(StepStatus::Running),
2351                    ..StepUpdate::default()
2352                },
2353            )
2354            .await
2355            .unwrap();
2356
2357        store
2358            .update_step(
2359                step.id,
2360                StepUpdate {
2361                    status: Some(StepStatus::Failed),
2362                    error: Some("Connection timeout after 30s".to_string()),
2363                    duration_ms: Some(30000),
2364                    ..StepUpdate::default()
2365                },
2366            )
2367            .await
2368            .unwrap();
2369
2370        let steps = store.list_steps(run.id).await.unwrap();
2371        assert_eq!(steps[0].status.state, StepStatus::Failed);
2372        assert_eq!(
2373            steps[0].error,
2374            Some("Connection timeout after 30s".to_string())
2375        );
2376        assert_eq!(steps[0].duration_ms, 30000);
2377    }
2378
2379    #[tokio::test]
2380    async fn list_steps_for_nonexistent_run_returns_empty() {
2381        let store = InMemoryStore::new();
2382        let steps = store.list_steps(Uuid::nil()).await.unwrap();
2383        assert!(steps.is_empty());
2384    }
2385
2386    #[tokio::test]
2387    async fn update_step_pending_to_skipped() {
2388        let store = InMemoryStore::new();
2389        let run = store
2390            .create_run(new_run_req("test"))
2391            .await
2392            .unwrap()
2393            .into_run();
2394
2395        let step = store
2396            .create_step(NewStep {
2397                run_id: run.id,
2398                trace_id: step_trace_id(run.id, "build", 0),
2399                name: "build".to_string(),
2400                kind: crate::entities::StepKind::Shell,
2401                position: 0,
2402                input: None,
2403                is_error_handler: false,
2404            })
2405            .await
2406            .unwrap();
2407
2408        // Pending → Skipped is allowed (when prior step failed)
2409        store
2410            .update_step(
2411                step.id,
2412                StepUpdate {
2413                    status: Some(StepStatus::Skipped),
2414                    ..StepUpdate::default()
2415                },
2416            )
2417            .await
2418            .unwrap();
2419
2420        let steps = store.list_steps(run.id).await.unwrap();
2421        assert_eq!(steps[0].status.state, StepStatus::Skipped);
2422    }
2423
2424    #[tokio::test]
2425    async fn list_runs_with_combined_filters() {
2426        let store = InMemoryStore::new();
2427
2428        let r1 = store
2429            .create_run(new_run_req("deploy"))
2430            .await
2431            .unwrap()
2432            .into_run();
2433        let r2 = store
2434            .create_run(new_run_req("deploy"))
2435            .await
2436            .unwrap()
2437            .into_run();
2438        let _r3 = store
2439            .create_run(new_run_req("test"))
2440            .await
2441            .unwrap()
2442            .into_run();
2443
2444        // r1: Pending → Running → Completed
2445        store
2446            .update_run_status(r1.id, RunStatus::Running)
2447            .await
2448            .unwrap();
2449        store
2450            .update_run_status(r1.id, RunStatus::Completed)
2451            .await
2452            .unwrap();
2453
2454        // r2: Pending → Running
2455        store
2456            .update_run_status(r2.id, RunStatus::Running)
2457            .await
2458            .unwrap();
2459
2460        // Filter by workflow AND status
2461        let filter = RunFilter {
2462            workflow_name: Some("deploy".to_string()),
2463            status: Some(RunStatus::Running),
2464            ..RunFilter::default()
2465        };
2466
2467        let page = store.list_runs(filter, 1, 100).await.unwrap();
2468        assert_eq!(page.total, 1);
2469        assert_eq!(page.items[0].id, r2.id);
2470    }
2471
2472    #[tokio::test]
2473    async fn list_runs_workflow_filter_is_case_insensitive_partial_match() {
2474        let store = InMemoryStore::new();
2475        store
2476            .create_run(new_run_req("weather-report"))
2477            .await
2478            .unwrap()
2479            .into_run();
2480        store
2481            .create_run(new_run_req("deploy-prod"))
2482            .await
2483            .unwrap()
2484            .into_run();
2485
2486        // Partial match
2487        let filter = RunFilter {
2488            workflow_name: Some("weather".to_string()),
2489            ..RunFilter::default()
2490        };
2491        let page = store.list_runs(filter, 1, 100).await.unwrap();
2492        assert_eq!(page.total, 1);
2493        assert_eq!(page.items[0].workflow_name, "weather-report");
2494
2495        // Case-insensitive
2496        let filter = RunFilter {
2497            workflow_name: Some("Weather-REPORT".to_string()),
2498            ..RunFilter::default()
2499        };
2500        let page = store.list_runs(filter, 1, 100).await.unwrap();
2501        assert_eq!(page.total, 1);
2502        assert_eq!(page.items[0].workflow_name, "weather-report");
2503
2504        // Suffix match
2505        let filter = RunFilter {
2506            workflow_name: Some("report".to_string()),
2507            ..RunFilter::default()
2508        };
2509        let page = store.list_runs(filter, 1, 100).await.unwrap();
2510        assert_eq!(page.total, 1);
2511        assert_eq!(page.items[0].workflow_name, "weather-report");
2512
2513        // No match
2514        let filter = RunFilter {
2515            workflow_name: Some("build".to_string()),
2516            ..RunFilter::default()
2517        };
2518        let page = store.list_runs(filter, 1, 100).await.unwrap();
2519        assert_eq!(page.total, 0);
2520    }
2521
2522    #[tokio::test]
2523    async fn list_runs_has_steps_true_only_filters_completed_and_cancelled() {
2524        let store = InMemoryStore::new();
2525        let run_with = create_terminal_run(&store, "with-steps", RunStatus::Completed).await;
2526        let _run_without = create_terminal_run(&store, "without-steps", RunStatus::Completed).await;
2527
2528        store
2529            .create_step(NewStep {
2530                run_id: run_with.id,
2531                trace_id: step_trace_id(run_with.id, "build", 0),
2532                name: "build".to_string(),
2533                kind: crate::entities::StepKind::Shell,
2534                position: 0,
2535                input: None,
2536                is_error_handler: false,
2537            })
2538            .await
2539            .unwrap();
2540
2541        let filter = RunFilter {
2542            has_steps: Some(true),
2543            ..RunFilter::default()
2544        };
2545        let page = store.list_runs(filter, 1, 100).await.unwrap();
2546        assert_eq!(page.total, 1);
2547        assert_eq!(page.items[0].id, run_with.id);
2548    }
2549
2550    #[tokio::test]
2551    async fn list_runs_has_steps_false_only_filters_completed_and_cancelled() {
2552        let store = InMemoryStore::new();
2553        let run_with = create_terminal_run(&store, "with-steps", RunStatus::Cancelled).await;
2554        let run_without = create_terminal_run(&store, "without-steps", RunStatus::Cancelled).await;
2555
2556        store
2557            .create_step(NewStep {
2558                run_id: run_with.id,
2559                trace_id: step_trace_id(run_with.id, "build", 0),
2560                name: "build".to_string(),
2561                kind: crate::entities::StepKind::Shell,
2562                position: 0,
2563                input: None,
2564                is_error_handler: false,
2565            })
2566            .await
2567            .unwrap();
2568
2569        let filter = RunFilter {
2570            has_steps: Some(false),
2571            ..RunFilter::default()
2572        };
2573        let page = store.list_runs(filter, 1, 100).await.unwrap();
2574        assert_eq!(page.total, 1);
2575        assert_eq!(page.items[0].id, run_without.id);
2576    }
2577
2578    #[tokio::test]
2579    async fn list_runs_has_steps_none_returns_all() {
2580        let store = InMemoryStore::new();
2581        let run_with = store
2582            .create_run(new_run_req("with-steps"))
2583            .await
2584            .unwrap()
2585            .into_run();
2586        let _run_without = store
2587            .create_run(new_run_req("without-steps"))
2588            .await
2589            .unwrap()
2590            .into_run();
2591
2592        store
2593            .create_step(NewStep {
2594                run_id: run_with.id,
2595                trace_id: step_trace_id(run_with.id, "build", 0),
2596                name: "build".to_string(),
2597                kind: crate::entities::StepKind::Shell,
2598                position: 0,
2599                input: None,
2600                is_error_handler: false,
2601            })
2602            .await
2603            .unwrap();
2604
2605        let filter = RunFilter {
2606            has_steps: None,
2607            ..RunFilter::default()
2608        };
2609        let page = store.list_runs(filter, 1, 100).await.unwrap();
2610        assert_eq!(page.total, 2);
2611    }
2612
2613    #[tokio::test]
2614    async fn list_runs_has_steps_true_does_not_filter_non_terminal_runs() {
2615        let store = InMemoryStore::new();
2616        let pending_run = store
2617            .create_run(new_run_req("pending-empty"))
2618            .await
2619            .unwrap()
2620            .into_run();
2621        let running_run = store
2622            .create_run(new_run_req("running-empty"))
2623            .await
2624            .unwrap()
2625            .into_run();
2626        store
2627            .update_run_status(running_run.id, RunStatus::Running)
2628            .await
2629            .unwrap();
2630
2631        let filter = RunFilter {
2632            has_steps: Some(true),
2633            ..RunFilter::default()
2634        };
2635        let page = store.list_runs(filter, 1, 100).await.unwrap();
2636        assert_eq!(page.total, 2);
2637        let ids: Vec<_> = page.items.iter().map(|r| r.id).collect();
2638        assert!(ids.contains(&pending_run.id));
2639        assert!(ids.contains(&running_run.id));
2640    }
2641
2642    #[tokio::test]
2643    async fn get_stats_with_mixed_active_statuses() {
2644        let store = InMemoryStore::new();
2645
2646        let _r1 = store
2647            .create_run(new_run_req("wf"))
2648            .await
2649            .unwrap()
2650            .into_run(); // Pending
2651        let r2 = store
2652            .create_run(new_run_req("wf"))
2653            .await
2654            .unwrap()
2655            .into_run();
2656        let r3 = store
2657            .create_run(new_run_req("wf"))
2658            .await
2659            .unwrap()
2660            .into_run();
2661
2662        store
2663            .update_run_status(r2.id, RunStatus::Running)
2664            .await
2665            .unwrap();
2666        store
2667            .update_run_status(r3.id, RunStatus::Running)
2668            .await
2669            .unwrap();
2670        store
2671            .update_run_status(r3.id, RunStatus::Retrying)
2672            .await
2673            .unwrap();
2674
2675        let r4 = store
2676            .create_run(new_run_req("wf"))
2677            .await
2678            .unwrap()
2679            .into_run();
2680        store
2681            .update_run_status(r4.id, RunStatus::Running)
2682            .await
2683            .unwrap();
2684        store
2685            .update_run_status(r4.id, RunStatus::AwaitingApproval)
2686            .await
2687            .unwrap();
2688
2689        let r5 = store
2690            .create_run(new_run_req("wf"))
2691            .await
2692            .unwrap()
2693            .into_run();
2694        store
2695            .update_run_status(r5.id, RunStatus::Running)
2696            .await
2697            .unwrap();
2698        store
2699            .update_run_status(r5.id, RunStatus::Sleeping)
2700            .await
2701            .unwrap();
2702
2703        let stats = store.get_stats(RunFilter::default()).await.unwrap();
2704        // _r1 (Pending), r2 (Running), r3 (Retrying), r4 (AwaitingApproval), r5 (Sleeping)
2705        assert_eq!(stats.active_runs, 5);
2706        assert_eq!(stats.awaiting_approval_runs, 1);
2707    }
2708
2709    #[tokio::test]
2710    async fn run_with_different_trigger_kinds() {
2711        let store = InMemoryStore::new();
2712
2713        let r1 = store
2714            .create_run(NewRun {
2715                created_by: None,
2716                workflow_name: "test".to_string(),
2717                trigger: TriggerKind::Manual,
2718                payload: json!({}),
2719                max_retries: 1,
2720                handler_version: None,
2721                labels: HashMap::new(),
2722                scheduled_at: None,
2723                idempotency_key: None,
2724                concurrency_key: None,
2725                concurrency_limits: Vec::new(),
2726                max_cost_usd: None,
2727                worker_tags: Vec::new(),
2728            })
2729            .await
2730            .unwrap()
2731            .into_run();
2732
2733        let r2 = store
2734            .create_run(NewRun {
2735                created_by: None,
2736                workflow_name: "test".to_string(),
2737                trigger: TriggerKind::Webhook {
2738                    path: "/hooks/github".to_string(),
2739                },
2740                payload: json!({}),
2741                max_retries: 1,
2742                handler_version: None,
2743                labels: HashMap::new(),
2744                scheduled_at: None,
2745                idempotency_key: None,
2746                concurrency_key: None,
2747                concurrency_limits: Vec::new(),
2748                max_cost_usd: None,
2749                worker_tags: Vec::new(),
2750            })
2751            .await
2752            .unwrap()
2753            .into_run();
2754
2755        let r3 = store
2756            .create_run(NewRun {
2757                created_by: None,
2758                workflow_name: "test".to_string(),
2759                trigger: TriggerKind::Cron {
2760                    schedule: "0 0 * * *".to_string(),
2761                },
2762                payload: json!({}),
2763                max_retries: 1,
2764                handler_version: None,
2765                labels: HashMap::new(),
2766                scheduled_at: None,
2767                idempotency_key: None,
2768                concurrency_key: None,
2769                concurrency_limits: Vec::new(),
2770                max_cost_usd: None,
2771                worker_tags: Vec::new(),
2772            })
2773            .await
2774            .unwrap()
2775            .into_run();
2776
2777        let r4 = store
2778            .create_run(NewRun {
2779                created_by: None,
2780                workflow_name: "test".to_string(),
2781                trigger: TriggerKind::Api,
2782                payload: json!({}),
2783                max_retries: 1,
2784                handler_version: None,
2785                labels: HashMap::new(),
2786                scheduled_at: None,
2787                idempotency_key: None,
2788                concurrency_key: None,
2789                concurrency_limits: Vec::new(),
2790                max_cost_usd: None,
2791                worker_tags: Vec::new(),
2792            })
2793            .await
2794            .unwrap()
2795            .into_run();
2796
2797        let r5 = store
2798            .create_run(NewRun {
2799                created_by: None,
2800                workflow_name: "test".to_string(),
2801                trigger: TriggerKind::Retry {
2802                    parent_run_id: Uuid::nil(),
2803                },
2804                payload: json!({}),
2805                max_retries: 1,
2806                handler_version: None,
2807                labels: HashMap::new(),
2808                scheduled_at: None,
2809                idempotency_key: None,
2810                concurrency_key: None,
2811                concurrency_limits: Vec::new(),
2812                max_cost_usd: None,
2813                worker_tags: Vec::new(),
2814            })
2815            .await
2816            .unwrap()
2817            .into_run();
2818
2819        assert_eq!(r1.trigger, TriggerKind::Manual);
2820        assert!(matches!(r2.trigger, TriggerKind::Webhook { .. }));
2821        assert!(matches!(r3.trigger, TriggerKind::Cron { .. }));
2822        assert_eq!(r4.trigger, TriggerKind::Api);
2823        assert!(matches!(r5.trigger, TriggerKind::Retry { .. }));
2824    }
2825
2826    // ---- create_step_dependencies ----
2827
2828    #[tokio::test]
2829    async fn create_step_dependencies_stores_dependencies() {
2830        let store = InMemoryStore::new();
2831        let run = store
2832            .create_run(new_run_req("test"))
2833            .await
2834            .unwrap()
2835            .into_run();
2836
2837        let step1 = store
2838            .create_step(NewStep {
2839                run_id: run.id,
2840                trace_id: step_trace_id(run.id, "step1", 0),
2841                name: "step1".to_string(),
2842                kind: crate::entities::StepKind::Shell,
2843                position: 0,
2844                input: None,
2845                is_error_handler: false,
2846            })
2847            .await
2848            .unwrap();
2849
2850        let step2 = store
2851            .create_step(NewStep {
2852                run_id: run.id,
2853                trace_id: step_trace_id(run.id, "step2", 1),
2854                name: "step2".to_string(),
2855                kind: crate::entities::StepKind::Shell,
2856                position: 1,
2857                input: None,
2858                is_error_handler: false,
2859            })
2860            .await
2861            .unwrap();
2862
2863        let result = store
2864            .create_step_dependencies(vec![NewStepDependency {
2865                step_id: step2.id,
2866                depends_on: step1.id,
2867            }])
2868            .await;
2869
2870        assert!(result.is_ok());
2871
2872        let deps = store.list_step_dependencies(run.id).await.unwrap();
2873        assert_eq!(deps.len(), 1);
2874        assert_eq!(deps[0].step_id, step2.id);
2875        assert_eq!(deps[0].depends_on, step1.id);
2876    }
2877
2878    #[tokio::test]
2879    async fn create_step_dependencies_duplicate_dependencies_are_idempotent() {
2880        let store = InMemoryStore::new();
2881        let run = store
2882            .create_run(new_run_req("test"))
2883            .await
2884            .unwrap()
2885            .into_run();
2886
2887        let step1 = store
2888            .create_step(NewStep {
2889                run_id: run.id,
2890                trace_id: step_trace_id(run.id, "step1", 0),
2891                name: "step1".to_string(),
2892                kind: crate::entities::StepKind::Shell,
2893                position: 0,
2894                input: None,
2895                is_error_handler: false,
2896            })
2897            .await
2898            .unwrap();
2899
2900        let step2 = store
2901            .create_step(NewStep {
2902                run_id: run.id,
2903                trace_id: step_trace_id(run.id, "step2", 1),
2904                name: "step2".to_string(),
2905                kind: crate::entities::StepKind::Shell,
2906                position: 1,
2907                input: None,
2908                is_error_handler: false,
2909            })
2910            .await
2911            .unwrap();
2912
2913        let dep = NewStepDependency {
2914            step_id: step2.id,
2915            depends_on: step1.id,
2916        };
2917
2918        store
2919            .create_step_dependencies(vec![dep.clone()])
2920            .await
2921            .unwrap();
2922        store.create_step_dependencies(vec![dep]).await.unwrap();
2923
2924        let deps = store.list_step_dependencies(run.id).await.unwrap();
2925        assert_eq!(deps.len(), 1);
2926    }
2927
2928    #[tokio::test]
2929    async fn create_step_dependencies_missing_step_id_returns_error() {
2930        let store = InMemoryStore::new();
2931        let run = store
2932            .create_run(new_run_req("test"))
2933            .await
2934            .unwrap()
2935            .into_run();
2936
2937        let step1 = store
2938            .create_step(NewStep {
2939                run_id: run.id,
2940                trace_id: step_trace_id(run.id, "step1", 0),
2941                name: "step1".to_string(),
2942                kind: crate::entities::StepKind::Shell,
2943                position: 0,
2944                input: None,
2945                is_error_handler: false,
2946            })
2947            .await
2948            .unwrap();
2949
2950        let result = store
2951            .create_step_dependencies(vec![NewStepDependency {
2952                step_id: Uuid::nil(),
2953                depends_on: step1.id,
2954            }])
2955            .await;
2956
2957        assert!(matches!(result.unwrap_err(), StoreError::StepNotFound(_)));
2958    }
2959
2960    #[tokio::test]
2961    async fn create_step_dependencies_missing_depends_on_returns_error() {
2962        let store = InMemoryStore::new();
2963        let run = store
2964            .create_run(new_run_req("test"))
2965            .await
2966            .unwrap()
2967            .into_run();
2968
2969        let step1 = store
2970            .create_step(NewStep {
2971                run_id: run.id,
2972                trace_id: step_trace_id(run.id, "step1", 0),
2973                name: "step1".to_string(),
2974                kind: crate::entities::StepKind::Shell,
2975                position: 0,
2976                input: None,
2977                is_error_handler: false,
2978            })
2979            .await
2980            .unwrap();
2981
2982        let result = store
2983            .create_step_dependencies(vec![NewStepDependency {
2984                step_id: step1.id,
2985                depends_on: Uuid::nil(),
2986            }])
2987            .await;
2988
2989        assert!(matches!(result.unwrap_err(), StoreError::StepNotFound(_)));
2990    }
2991
2992    #[tokio::test]
2993    async fn create_step_dependencies_multiple_dependencies() {
2994        let store = InMemoryStore::new();
2995        let run = store
2996            .create_run(new_run_req("test"))
2997            .await
2998            .unwrap()
2999            .into_run();
3000
3001        let step1 = store
3002            .create_step(NewStep {
3003                run_id: run.id,
3004                trace_id: step_trace_id(run.id, "step1", 0),
3005                name: "step1".to_string(),
3006                kind: crate::entities::StepKind::Shell,
3007                position: 0,
3008                input: None,
3009                is_error_handler: false,
3010            })
3011            .await
3012            .unwrap();
3013
3014        let step2 = store
3015            .create_step(NewStep {
3016                run_id: run.id,
3017                trace_id: step_trace_id(run.id, "step2", 1),
3018                name: "step2".to_string(),
3019                kind: crate::entities::StepKind::Shell,
3020                position: 1,
3021                input: None,
3022                is_error_handler: false,
3023            })
3024            .await
3025            .unwrap();
3026
3027        let step3 = store
3028            .create_step(NewStep {
3029                run_id: run.id,
3030                trace_id: step_trace_id(run.id, "step3", 2),
3031                name: "step3".to_string(),
3032                kind: crate::entities::StepKind::Shell,
3033                position: 2,
3034                input: None,
3035                is_error_handler: false,
3036            })
3037            .await
3038            .unwrap();
3039
3040        let result = store
3041            .create_step_dependencies(vec![
3042                NewStepDependency {
3043                    step_id: step2.id,
3044                    depends_on: step1.id,
3045                },
3046                NewStepDependency {
3047                    step_id: step3.id,
3048                    depends_on: step2.id,
3049                },
3050            ])
3051            .await;
3052
3053        assert!(result.is_ok());
3054
3055        let deps = store.list_step_dependencies(run.id).await.unwrap();
3056        assert_eq!(deps.len(), 2);
3057    }
3058
3059    // ---- list_step_dependencies ----
3060
3061    #[tokio::test]
3062    async fn list_step_dependencies_empty_for_run_with_no_dependencies() {
3063        let store = InMemoryStore::new();
3064        let run = store
3065            .create_run(new_run_req("test"))
3066            .await
3067            .unwrap()
3068            .into_run();
3069
3070        store
3071            .create_step(NewStep {
3072                run_id: run.id,
3073                trace_id: step_trace_id(run.id, "step1", 0),
3074                name: "step1".to_string(),
3075                kind: crate::entities::StepKind::Shell,
3076                position: 0,
3077                input: None,
3078                is_error_handler: false,
3079            })
3080            .await
3081            .unwrap();
3082
3083        let deps = store.list_step_dependencies(run.id).await.unwrap();
3084        assert!(deps.is_empty());
3085    }
3086
3087    #[tokio::test]
3088    async fn list_step_dependencies_returns_only_deps_for_given_run() {
3089        let store = InMemoryStore::new();
3090        let run1 = store
3091            .create_run(new_run_req("test1"))
3092            .await
3093            .unwrap()
3094            .into_run();
3095        let run2 = store
3096            .create_run(new_run_req("test2"))
3097            .await
3098            .unwrap()
3099            .into_run();
3100
3101        let step1_run1 = store
3102            .create_step(NewStep {
3103                run_id: run1.id,
3104                trace_id: step_trace_id(run1.id, "step1", 0),
3105                name: "step1".to_string(),
3106                kind: crate::entities::StepKind::Shell,
3107                position: 0,
3108                input: None,
3109                is_error_handler: false,
3110            })
3111            .await
3112            .unwrap();
3113
3114        let step2_run1 = store
3115            .create_step(NewStep {
3116                run_id: run1.id,
3117                trace_id: step_trace_id(run1.id, "step2", 1),
3118                name: "step2".to_string(),
3119                kind: crate::entities::StepKind::Shell,
3120                position: 1,
3121                input: None,
3122                is_error_handler: false,
3123            })
3124            .await
3125            .unwrap();
3126
3127        let step1_run2 = store
3128            .create_step(NewStep {
3129                run_id: run2.id,
3130                trace_id: step_trace_id(run2.id, "step1", 0),
3131                name: "step1".to_string(),
3132                kind: crate::entities::StepKind::Shell,
3133                position: 0,
3134                input: None,
3135                is_error_handler: false,
3136            })
3137            .await
3138            .unwrap();
3139
3140        let step2_run2 = store
3141            .create_step(NewStep {
3142                run_id: run2.id,
3143                trace_id: step_trace_id(run2.id, "step2", 1),
3144                name: "step2".to_string(),
3145                kind: crate::entities::StepKind::Shell,
3146                position: 1,
3147                input: None,
3148                is_error_handler: false,
3149            })
3150            .await
3151            .unwrap();
3152
3153        store
3154            .create_step_dependencies(vec![
3155                NewStepDependency {
3156                    step_id: step2_run1.id,
3157                    depends_on: step1_run1.id,
3158                },
3159                NewStepDependency {
3160                    step_id: step2_run2.id,
3161                    depends_on: step1_run2.id,
3162                },
3163            ])
3164            .await
3165            .unwrap();
3166
3167        let deps_run1 = store.list_step_dependencies(run1.id).await.unwrap();
3168        let deps_run2 = store.list_step_dependencies(run2.id).await.unwrap();
3169
3170        assert_eq!(deps_run1.len(), 1);
3171        assert_eq!(deps_run1[0].step_id, step2_run1.id);
3172        assert_eq!(deps_run1[0].depends_on, step1_run1.id);
3173
3174        assert_eq!(deps_run2.len(), 1);
3175        assert_eq!(deps_run2[0].step_id, step2_run2.id);
3176        assert_eq!(deps_run2[0].depends_on, step1_run2.id);
3177    }
3178
3179    #[tokio::test]
3180    async fn list_step_dependencies_returns_empty_for_nonexistent_run() {
3181        let store = InMemoryStore::new();
3182        let deps = store.list_step_dependencies(Uuid::nil()).await.unwrap();
3183        assert!(deps.is_empty());
3184    }
3185
3186    #[tokio::test]
3187    async fn list_step_dependencies_sorted_by_created_at() {
3188        let store = InMemoryStore::new();
3189        let run = store
3190            .create_run(new_run_req("test"))
3191            .await
3192            .unwrap()
3193            .into_run();
3194
3195        let step1 = store
3196            .create_step(NewStep {
3197                run_id: run.id,
3198                trace_id: step_trace_id(run.id, "step1", 0),
3199                name: "step1".to_string(),
3200                kind: crate::entities::StepKind::Shell,
3201                position: 0,
3202                input: None,
3203                is_error_handler: false,
3204            })
3205            .await
3206            .unwrap();
3207
3208        let step2 = store
3209            .create_step(NewStep {
3210                run_id: run.id,
3211                trace_id: step_trace_id(run.id, "step2", 1),
3212                name: "step2".to_string(),
3213                kind: crate::entities::StepKind::Shell,
3214                position: 1,
3215                input: None,
3216                is_error_handler: false,
3217            })
3218            .await
3219            .unwrap();
3220
3221        let step3 = store
3222            .create_step(NewStep {
3223                run_id: run.id,
3224                trace_id: step_trace_id(run.id, "step3", 2),
3225                name: "step3".to_string(),
3226                kind: crate::entities::StepKind::Shell,
3227                position: 2,
3228                input: None,
3229                is_error_handler: false,
3230            })
3231            .await
3232            .unwrap();
3233
3234        store
3235            .create_step_dependencies(vec![NewStepDependency {
3236                step_id: step2.id,
3237                depends_on: step1.id,
3238            }])
3239            .await
3240            .unwrap();
3241
3242        store
3243            .create_step_dependencies(vec![NewStepDependency {
3244                step_id: step3.id,
3245                depends_on: step1.id,
3246            }])
3247            .await
3248            .unwrap();
3249
3250        let deps = store.list_step_dependencies(run.id).await.unwrap();
3251        assert_eq!(deps.len(), 2);
3252        assert!(deps[0].created_at <= deps[1].created_at);
3253    }
3254
3255    // ---- update_run_returning ----
3256
3257    #[tokio::test]
3258    async fn update_run_returning_applies_and_returns() {
3259        let store = InMemoryStore::new();
3260        let run = store
3261            .create_run(new_run_req("test"))
3262            .await
3263            .unwrap()
3264            .into_run();
3265
3266        // Transition Pending -> Running first
3267        store
3268            .update_run_status(run.id, RunStatus::Running)
3269            .await
3270            .unwrap();
3271
3272        let updated = store
3273            .update_run_returning(
3274                run.id,
3275                RunUpdate {
3276                    status: Some(RunStatus::Completed),
3277                    cost_usd: Some(Decimal::new(4200, 2)),
3278                    duration_ms: Some(1500),
3279                    ..RunUpdate::default()
3280                },
3281            )
3282            .await
3283            .unwrap();
3284
3285        assert_eq!(updated.id, run.id);
3286        assert_eq!(updated.status.state, RunStatus::Completed);
3287        assert_eq!(updated.cost_usd, Decimal::new(4200, 2));
3288        assert_eq!(updated.duration_ms, 1500);
3289        assert!(updated.completed_at.is_some());
3290    }
3291
3292    #[tokio::test]
3293    async fn update_run_returning_not_found() {
3294        let store = InMemoryStore::new();
3295        let result = store
3296            .update_run_returning(
3297                Uuid::nil(),
3298                RunUpdate {
3299                    status: Some(RunStatus::Running),
3300                    ..RunUpdate::default()
3301                },
3302            )
3303            .await;
3304
3305        assert!(matches!(result, Err(StoreError::RunNotFound(_))));
3306    }
3307
3308    #[tokio::test]
3309    async fn update_run_returning_invalid_transition() {
3310        let store = InMemoryStore::new();
3311        let run = store
3312            .create_run(new_run_req("test"))
3313            .await
3314            .unwrap()
3315            .into_run();
3316
3317        let result = store
3318            .update_run_returning(
3319                run.id,
3320                RunUpdate {
3321                    status: Some(RunStatus::Completed),
3322                    ..RunUpdate::default()
3323                },
3324            )
3325            .await;
3326
3327        assert!(matches!(result, Err(StoreError::InvalidTransition { .. })));
3328    }
3329
3330    // ---- retry scheduling ----
3331
3332    #[tokio::test]
3333    async fn create_step_stamps_the_current_attempt() {
3334        let store = InMemoryStore::new();
3335        let run = store
3336            .create_run(new_run_req("retry-wf"))
3337            .await
3338            .unwrap()
3339            .into_run();
3340
3341        let first = store
3342            .create_step(new_step_req(run.id, "build", 0))
3343            .await
3344            .unwrap();
3345        assert_eq!(first.attempt, 1);
3346
3347        store
3348            .update_run_status(run.id, RunStatus::Running)
3349            .await
3350            .unwrap();
3351        store
3352            .update_run(
3353                run.id,
3354                RunUpdate {
3355                    status: Some(RunStatus::Retrying),
3356                    increment_retry: true,
3357                    ..RunUpdate::default()
3358                },
3359            )
3360            .await
3361            .unwrap();
3362
3363        let second = store
3364            .create_step(new_step_req(run.id, "build", 0))
3365            .await
3366            .unwrap();
3367        assert_eq!(second.attempt, 2);
3368    }
3369
3370    #[tokio::test]
3371    async fn pick_next_pending_ignores_retrying_run_before_its_backoff() {
3372        let store = InMemoryStore::new();
3373        let run = store
3374            .create_run(new_run_req("retry-wf"))
3375            .await
3376            .unwrap()
3377            .into_run();
3378
3379        store
3380            .update_run_status(run.id, RunStatus::Running)
3381            .await
3382            .unwrap();
3383        store
3384            .update_run(
3385                run.id,
3386                RunUpdate {
3387                    status: Some(RunStatus::Retrying),
3388                    increment_retry: true,
3389                    scheduled_at: Some(Utc::now() + TimeDelta::seconds(60)),
3390                    ..RunUpdate::default()
3391                },
3392            )
3393            .await
3394            .unwrap();
3395
3396        assert!(store.pick_next_pending(None).await.unwrap().is_none());
3397    }
3398
3399    #[tokio::test]
3400    async fn pick_next_pending_resumes_retrying_run_after_its_backoff() {
3401        let store = InMemoryStore::new();
3402        let run = store
3403            .create_run(new_run_req("retry-wf"))
3404            .await
3405            .unwrap()
3406            .into_run();
3407
3408        store
3409            .update_run_status(run.id, RunStatus::Running)
3410            .await
3411            .unwrap();
3412        store
3413            .update_run(
3414                run.id,
3415                RunUpdate {
3416                    status: Some(RunStatus::Retrying),
3417                    increment_retry: true,
3418                    scheduled_at: Some(Utc::now() - TimeDelta::seconds(1)),
3419                    ..RunUpdate::default()
3420                },
3421            )
3422            .await
3423            .unwrap();
3424
3425        let picked = store.pick_next_pending(None).await.unwrap().unwrap();
3426        assert_eq!(picked.id, run.id);
3427        assert_eq!(picked.status.state, RunStatus::Running);
3428        assert_eq!(picked.retry_count, 1);
3429    }
3430
3431    #[tokio::test]
3432    async fn update_run_persists_scheduled_at() {
3433        let store = InMemoryStore::new();
3434        let run = store
3435            .create_run(new_run_req("test"))
3436            .await
3437            .unwrap()
3438            .into_run();
3439        let when = Utc::now() + TimeDelta::seconds(30);
3440
3441        store
3442            .update_run(
3443                run.id,
3444                RunUpdate {
3445                    scheduled_at: Some(when),
3446                    ..RunUpdate::default()
3447                },
3448            )
3449            .await
3450            .unwrap();
3451
3452        let fetched = store.get_run(run.id).await.unwrap().unwrap();
3453        assert_eq!(fetched.scheduled_at, Some(when));
3454    }
3455
3456    // ---- output ----
3457
3458    #[tokio::test]
3459    async fn new_run_has_no_output() {
3460        let store = InMemoryStore::new();
3461        let run = store
3462            .create_run(new_run_req("test"))
3463            .await
3464            .unwrap()
3465            .into_run();
3466
3467        assert!(run.output.is_none());
3468        let fetched = store.get_run(run.id).await.unwrap().unwrap();
3469        assert!(fetched.output.is_none());
3470    }
3471
3472    #[tokio::test]
3473    async fn update_run_sets_output() {
3474        let store = InMemoryStore::new();
3475        let run = store
3476            .create_run(new_run_req("test"))
3477            .await
3478            .unwrap()
3479            .into_run();
3480
3481        store
3482            .update_run(
3483                run.id,
3484                RunUpdate {
3485                    output: Some(json!({"verdict": "approved"})),
3486                    ..RunUpdate::default()
3487                },
3488            )
3489            .await
3490            .unwrap();
3491
3492        let fetched = store.get_run(run.id).await.unwrap().unwrap();
3493        assert_eq!(fetched.output, Some(json!({"verdict": "approved"})));
3494    }
3495
3496    #[tokio::test]
3497    async fn update_run_without_output_keeps_previous_output() {
3498        let store = InMemoryStore::new();
3499        let run = store
3500            .create_run(new_run_req("test"))
3501            .await
3502            .unwrap()
3503            .into_run();
3504
3505        store
3506            .update_run(
3507                run.id,
3508                RunUpdate {
3509                    output: Some(json!({"verdict": "approved"})),
3510                    ..RunUpdate::default()
3511                },
3512            )
3513            .await
3514            .unwrap();
3515        store
3516            .update_run(
3517                run.id,
3518                RunUpdate {
3519                    error: Some("boom".to_string()),
3520                    ..RunUpdate::default()
3521                },
3522            )
3523            .await
3524            .unwrap();
3525
3526        let fetched = store.get_run(run.id).await.unwrap().unwrap();
3527        assert_eq!(fetched.output, Some(json!({"verdict": "approved"})));
3528        assert_eq!(fetched.error.as_deref(), Some("boom"));
3529    }
3530
3531    #[tokio::test]
3532    async fn update_run_output_last_write_wins() {
3533        let store = InMemoryStore::new();
3534        let run = store
3535            .create_run(new_run_req("test"))
3536            .await
3537            .unwrap()
3538            .into_run();
3539
3540        for verdict in ["first", "second"] {
3541            store
3542                .update_run(
3543                    run.id,
3544                    RunUpdate {
3545                        output: Some(json!({ "verdict": verdict })),
3546                        ..RunUpdate::default()
3547                    },
3548                )
3549                .await
3550                .unwrap();
3551        }
3552
3553        let fetched = store.get_run(run.id).await.unwrap().unwrap();
3554        assert_eq!(fetched.output, Some(json!({"verdict": "second"})));
3555    }
3556
3557    // ---- created_by ----
3558
3559    async fn seed_user(store: &InMemoryStore, username: &str) -> Uuid {
3560        store
3561            .create_user(NewUser {
3562                email: format!("{username}@example.com"),
3563                username: username.to_string(),
3564                password_hash: "hash".to_string(),
3565                is_admin: Some(false),
3566            })
3567            .await
3568            .unwrap()
3569            .id
3570    }
3571
3572    async fn seed_api_key(store: &InMemoryStore, user_id: Uuid, name: &str) -> Uuid {
3573        store
3574            .create_api_key(NewApiKey {
3575                user_id,
3576                name: name.to_string(),
3577                key_hash: "hash".to_string(),
3578                key_prefix: "irfl_0000".to_string(),
3579                scopes: vec![ApiKeyScope::RunsWrite],
3580                expires_at: None,
3581                rate_limit_override: None,
3582            })
3583            .await
3584            .unwrap()
3585            .id
3586    }
3587
3588    fn run_req_by(actor: RunActor) -> NewRun {
3589        NewRun {
3590            created_by: Some(actor),
3591            ..new_run_req("test")
3592        }
3593    }
3594
3595    #[tokio::test]
3596    async fn create_run_without_actor_has_no_author() {
3597        let store = InMemoryStore::new();
3598        let run = store
3599            .create_run(new_run_req("test"))
3600            .await
3601            .unwrap()
3602            .into_run();
3603
3604        assert!(run.created_by.is_none());
3605        assert!(run.created_by_label.is_none());
3606    }
3607
3608    #[tokio::test]
3609    async fn create_run_by_user_resolves_username_as_label() {
3610        let store = InMemoryStore::new();
3611        let user_id = seed_user(&store, "alice").await;
3612
3613        let run = store
3614            .create_run(run_req_by(RunActor::User { user_id }))
3615            .await
3616            .unwrap()
3617            .into_run();
3618
3619        assert_eq!(run.created_by, Some(RunActor::User { user_id }));
3620        assert_eq!(run.created_by_label.as_deref(), Some("alice"));
3621    }
3622
3623    #[tokio::test]
3624    async fn create_run_by_api_key_resolves_key_and_owner_as_label() {
3625        let store = InMemoryStore::new();
3626        let user_id = seed_user(&store, "alice").await;
3627        let api_key_id = seed_api_key(&store, user_id, "ci-deploy").await;
3628
3629        let run = store
3630            .create_run(run_req_by(RunActor::ApiKey {
3631                api_key_id,
3632                user_id,
3633            }))
3634            .await
3635            .unwrap()
3636            .into_run();
3637
3638        assert_eq!(run.created_by_label.as_deref(), Some("ci-deploy (alice)"));
3639    }
3640
3641    #[tokio::test]
3642    async fn label_follows_api_key_rename() {
3643        let store = InMemoryStore::new();
3644        let user_id = seed_user(&store, "alice").await;
3645        let api_key_id = seed_api_key(&store, user_id, "ci-deploy").await;
3646        let run = store
3647            .create_run(run_req_by(RunActor::ApiKey {
3648                api_key_id,
3649                user_id,
3650            }))
3651            .await
3652            .unwrap()
3653            .into_run();
3654
3655        store
3656            .update_api_key(
3657                api_key_id,
3658                ApiKeyUpdate {
3659                    name: Some("ci-release".to_string()),
3660                    ..ApiKeyUpdate::default()
3661                },
3662            )
3663            .await
3664            .unwrap();
3665
3666        let reread = store.get_run(run.id).await.unwrap().unwrap();
3667        assert_eq!(
3668            reread.created_by_label.as_deref(),
3669            Some("ci-release (alice)")
3670        );
3671    }
3672
3673    #[tokio::test]
3674    async fn label_is_none_when_user_is_unknown() {
3675        let store = InMemoryStore::new();
3676        let run = store
3677            .create_run(run_req_by(RunActor::User {
3678                user_id: Uuid::now_v7(),
3679            }))
3680            .await
3681            .unwrap()
3682            .into_run();
3683
3684        assert!(run.created_by.is_some());
3685        assert!(run.created_by_label.is_none());
3686    }
3687
3688    #[tokio::test]
3689    async fn label_is_key_name_only_when_owner_is_unknown() {
3690        let store = InMemoryStore::new();
3691        let owner = seed_user(&store, "alice").await;
3692        let api_key_id = seed_api_key(&store, owner, "ci-deploy").await;
3693
3694        // An actor pointing at a user that was never created.
3695        let run = store
3696            .create_run(run_req_by(RunActor::ApiKey {
3697                api_key_id,
3698                user_id: Uuid::now_v7(),
3699            }))
3700            .await
3701            .unwrap()
3702            .into_run();
3703
3704        assert_eq!(run.created_by_label.as_deref(), Some("ci-deploy"));
3705    }
3706
3707    #[tokio::test]
3708    async fn list_runs_filters_by_author() {
3709        let store = InMemoryStore::new();
3710        let alice = seed_user(&store, "alice").await;
3711        let bob = seed_user(&store, "bob").await;
3712
3713        store
3714            .create_run(run_req_by(RunActor::User { user_id: alice }))
3715            .await
3716            .unwrap()
3717            .into_run();
3718        store
3719            .create_run(run_req_by(RunActor::User { user_id: bob }))
3720            .await
3721            .unwrap()
3722            .into_run();
3723        store.create_run(new_run_req("anonymous")).await.unwrap();
3724
3725        let page = store
3726            .list_runs(
3727                RunFilter {
3728                    created_by_user_id: Some(alice),
3729                    ..RunFilter::default()
3730                },
3731                1,
3732                20,
3733            )
3734            .await
3735            .unwrap();
3736
3737        assert_eq!(page.total, 1);
3738        assert_eq!(page.items[0].created_by_label.as_deref(), Some("alice"));
3739    }
3740
3741    #[tokio::test]
3742    async fn list_runs_author_filter_matches_runs_from_the_users_api_keys() {
3743        let store = InMemoryStore::new();
3744        let alice = seed_user(&store, "alice").await;
3745        let api_key_id = seed_api_key(&store, alice, "ci-deploy").await;
3746
3747        store
3748            .create_run(run_req_by(RunActor::ApiKey {
3749                api_key_id,
3750                user_id: alice,
3751            }))
3752            .await
3753            .unwrap()
3754            .into_run();
3755
3756        let page = store
3757            .list_runs(
3758                RunFilter {
3759                    created_by_user_id: Some(alice),
3760                    ..RunFilter::default()
3761                },
3762                1,
3763                20,
3764            )
3765            .await
3766            .unwrap();
3767
3768        assert_eq!(page.total, 1);
3769    }
3770
3771    #[tokio::test]
3772    async fn list_runs_author_filter_excludes_unrelated_users() {
3773        let store = InMemoryStore::new();
3774        let alice = seed_user(&store, "alice").await;
3775
3776        store
3777            .create_run(run_req_by(RunActor::User { user_id: alice }))
3778            .await
3779            .unwrap()
3780            .into_run();
3781
3782        let page = store
3783            .list_runs(
3784                RunFilter {
3785                    created_by_user_id: Some(Uuid::now_v7()),
3786                    ..RunFilter::default()
3787                },
3788                1,
3789                20,
3790            )
3791            .await
3792            .unwrap();
3793
3794        assert_eq!(page.total, 0);
3795    }
3796
3797    #[tokio::test]
3798    async fn list_runs_without_author_filter_returns_every_run() {
3799        let store = InMemoryStore::new();
3800        let alice = seed_user(&store, "alice").await;
3801
3802        store
3803            .create_run(run_req_by(RunActor::User { user_id: alice }))
3804            .await
3805            .unwrap()
3806            .into_run();
3807        store.create_run(new_run_req("anonymous")).await.unwrap();
3808
3809        let page = store.list_runs(RunFilter::default(), 1, 20).await.unwrap();
3810        assert_eq!(page.total, 2);
3811    }
3812
3813    #[tokio::test]
3814    async fn pick_next_pending_resolves_author_label() {
3815        let store = InMemoryStore::new();
3816        let user_id = seed_user(&store, "alice").await;
3817        store
3818            .create_run(run_req_by(RunActor::User { user_id }))
3819            .await
3820            .unwrap()
3821            .into_run();
3822
3823        let picked = store.pick_next_pending(None).await.unwrap().unwrap();
3824        assert_eq!(picked.created_by_label.as_deref(), Some("alice"));
3825    }
3826
3827    // ---- list_purgeable_runs ----
3828
3829    #[tokio::test]
3830    async fn list_purgeable_runs_returns_old_terminal_runs() {
3831        let store = InMemoryStore::new();
3832        let old = create_terminal_run(&store, "deploy", RunStatus::Completed).await;
3833        store
3834            .set_run_created_at(old.id, Utc::now() - chrono::Duration::days(100))
3835            .await;
3836
3837        let policy = PurgePolicy {
3838            max_age_days: 90,
3839            max_runs_per_workflow: 10000,
3840            dry_run: false,
3841        };
3842        let result = store.list_purgeable_runs(&policy, 100).await.unwrap();
3843
3844        assert_eq!(result.len(), 1);
3845        assert_eq!(result[0].run_id, old.id);
3846        assert_eq!(result[0].reason, PurgeReason::TooOld);
3847    }
3848
3849    #[tokio::test]
3850    async fn list_purgeable_runs_ignores_non_terminal_states() {
3851        let store = InMemoryStore::new();
3852
3853        // Pending
3854        let pending = store
3855            .create_run(new_run_req("deploy"))
3856            .await
3857            .unwrap()
3858            .into_run();
3859        store
3860            .set_run_created_at(pending.id, Utc::now() - chrono::Duration::days(200))
3861            .await;
3862
3863        // Running
3864        let running = store
3865            .create_run(new_run_req("deploy"))
3866            .await
3867            .unwrap()
3868            .into_run();
3869        store
3870            .update_run_status(running.id, RunStatus::Running)
3871            .await
3872            .unwrap();
3873        store
3874            .set_run_created_at(running.id, Utc::now() - chrono::Duration::days(200))
3875            .await;
3876
3877        let policy = PurgePolicy {
3878            max_age_days: 90,
3879            max_runs_per_workflow: 1,
3880            dry_run: false,
3881        };
3882        let result = store.list_purgeable_runs(&policy, 100).await.unwrap();
3883        assert!(result.is_empty());
3884    }
3885
3886    #[tokio::test]
3887    async fn list_purgeable_runs_returns_excess_per_workflow() {
3888        let store = InMemoryStore::new();
3889        let r1 = create_terminal_run(&store, "deploy", RunStatus::Completed).await;
3890        store
3891            .set_run_created_at(r1.id, Utc::now() - chrono::Duration::days(10))
3892            .await;
3893        let r2 = create_terminal_run(&store, "deploy", RunStatus::Completed).await;
3894        store
3895            .set_run_created_at(r2.id, Utc::now() - chrono::Duration::days(5))
3896            .await;
3897        let _r3 = create_terminal_run(&store, "deploy", RunStatus::Completed).await;
3898
3899        let policy = PurgePolicy {
3900            max_age_days: 365,
3901            max_runs_per_workflow: 2,
3902            dry_run: false,
3903        };
3904        let result = store.list_purgeable_runs(&policy, 100).await.unwrap();
3905
3906        assert_eq!(result.len(), 1);
3907        assert_eq!(result[0].run_id, r1.id);
3908        assert_eq!(result[0].reason, PurgeReason::ExceedsWorkflowLimit);
3909    }
3910
3911    // ---- delete_run ----
3912
3913    #[tokio::test]
3914    async fn delete_run_removes_run_and_associated_data() {
3915        use crate::artifact_store::ArtifactStore;
3916        use crate::entities::{NewStep, StepKind};
3917
3918        let store = InMemoryStore::new();
3919        let run = create_terminal_run(&store, "deploy", RunStatus::Completed).await;
3920        let step = store
3921            .create_step(NewStep {
3922                run_id: run.id,
3923                trace_id: step_trace_id(run.id, "build", 0),
3924                name: "build".to_string(),
3925                kind: StepKind::Shell,
3926                position: 0,
3927                input: None,
3928                is_error_handler: false,
3929            })
3930            .await
3931            .unwrap();
3932
3933        let artifact_id = Uuid::now_v7();
3934        store
3935            .create_artifact(crate::entities::NewArtifact {
3936                id: artifact_id,
3937                run_id: run.id,
3938                step_id: step.id,
3939                name: "report.html".to_string(),
3940                storage_key: format!("artifacts/{}/{}/{}", run.id, step.id, artifact_id),
3941                content_type: "text/html".to_string(),
3942                size_bytes: 42,
3943                sha256: "0".repeat(64),
3944            })
3945            .await
3946            .unwrap();
3947
3948        let keys = store.delete_run(run.id).await.unwrap();
3949
3950        assert_eq!(keys.len(), 1);
3951        assert!(keys[0].contains(&artifact_id.to_string()));
3952        assert!(store.get_run(run.id).await.unwrap().is_none());
3953        assert!(store.list_steps(run.id).await.unwrap().is_empty());
3954        assert!(
3955            store
3956                .list_artifacts_for_run(run.id)
3957                .await
3958                .unwrap()
3959                .is_empty()
3960        );
3961    }
3962
3963    #[tokio::test]
3964    async fn delete_run_not_found() {
3965        let store = InMemoryStore::new();
3966        let err = store.delete_run(Uuid::now_v7()).await.unwrap_err();
3967        assert!(matches!(err, StoreError::RunNotFound(_)));
3968    }
3969
3970    // ---- record_step_approval ----
3971
3972    fn vote(user_id: Uuid, name: &str) -> StepApproval {
3973        StepApproval {
3974            user_id,
3975            approved_by: name.to_string(),
3976            at: Utc::now(),
3977        }
3978    }
3979
3980    #[tokio::test]
3981    async fn record_step_approval_appends_distinct_voters() {
3982        let store = InMemoryStore::new();
3983        let run = store
3984            .create_run(new_run_req("test"))
3985            .await
3986            .unwrap()
3987            .into_run();
3988        let step = store
3989            .create_step(new_step_req(run.id, "gate", 0))
3990            .await
3991            .unwrap();
3992        assert!(step.approvals.is_empty());
3993        assert!(step.approval_requirement.is_none());
3994
3995        let alice = Uuid::now_v7();
3996        let bob = Uuid::now_v7();
3997        let after_first = store
3998            .record_step_approval(step.id, vote(alice, "alice"))
3999            .await
4000            .unwrap();
4001        assert_eq!(after_first.approvals.len(), 1);
4002
4003        let after_second = store
4004            .record_step_approval(step.id, vote(bob, "bob"))
4005            .await
4006            .unwrap();
4007        assert_eq!(after_second.approvals.len(), 2);
4008        assert_eq!(after_second.approvals[0].user_id, alice);
4009        assert_eq!(after_second.approvals[1].user_id, bob);
4010    }
4011
4012    #[tokio::test]
4013    async fn record_step_approval_ignores_same_user() {
4014        let store = InMemoryStore::new();
4015        let run = store
4016            .create_run(new_run_req("test"))
4017            .await
4018            .unwrap()
4019            .into_run();
4020        let step = store
4021            .create_step(new_step_req(run.id, "gate", 0))
4022            .await
4023            .unwrap();
4024
4025        let alice = Uuid::now_v7();
4026        store
4027            .record_step_approval(step.id, vote(alice, "alice"))
4028            .await
4029            .unwrap();
4030        let again = store
4031            .record_step_approval(step.id, vote(alice, "alice-key"))
4032            .await
4033            .unwrap();
4034
4035        assert_eq!(again.approvals.len(), 1);
4036        assert_eq!(again.approvals[0].approved_by, "alice");
4037    }
4038
4039    #[tokio::test]
4040    async fn record_step_approval_unknown_step_is_not_found() {
4041        let store = InMemoryStore::new();
4042        let err = store
4043            .record_step_approval(Uuid::now_v7(), vote(Uuid::now_v7(), "alice"))
4044            .await
4045            .unwrap_err();
4046        assert!(matches!(err, StoreError::StepNotFound(_)));
4047    }
4048
4049    #[tokio::test]
4050    async fn update_step_sets_approval_requirement() {
4051        let store = InMemoryStore::new();
4052        let run = store
4053            .create_run(new_run_req("test"))
4054            .await
4055            .unwrap()
4056            .into_run();
4057        let step = store
4058            .create_step(new_step_req(run.id, "gate", 0))
4059            .await
4060            .unwrap();
4061        let requirement = ApprovalRequirement {
4062            required_approvers: 3,
4063            ..ApprovalRequirement::default()
4064        };
4065
4066        store
4067            .update_step(
4068                step.id,
4069                StepUpdate {
4070                    approval_requirement: Some(requirement.clone()),
4071                    ..StepUpdate::default()
4072                },
4073            )
4074            .await
4075            .unwrap();
4076
4077        let fetched = store.get_step(step.id).await.unwrap().unwrap();
4078        assert_eq!(fetched.approval_requirement, Some(requirement));
4079    }
4080
4081    async fn sleeping_run(store: &InMemoryStore, scheduled_at: DateTime<Utc>) -> Run {
4082        let run = store
4083            .create_run(new_run_req("sleepy"))
4084            .await
4085            .unwrap()
4086            .into_run();
4087        store
4088            .update_run_status(run.id, RunStatus::Running)
4089            .await
4090            .unwrap();
4091        store
4092            .update_run(
4093                run.id,
4094                RunUpdate {
4095                    status: Some(RunStatus::Sleeping),
4096                    scheduled_at: Some(scheduled_at),
4097                    ..RunUpdate::default()
4098                },
4099            )
4100            .await
4101            .unwrap();
4102        store.get_run(run.id).await.unwrap().unwrap()
4103    }
4104
4105    #[tokio::test]
4106    async fn claim_due_sleeping_runs_requeues_due_runs() {
4107        let store = InMemoryStore::new();
4108        let due = sleeping_run(&store, Utc::now() - TimeDelta::seconds(5)).await;
4109
4110        let woken = store.claim_due_sleeping_runs(10).await.unwrap();
4111        assert_eq!(woken.len(), 1);
4112        assert_eq!(woken[0].id, due.id);
4113        assert_eq!(woken[0].status.state, RunStatus::Pending);
4114        assert!(woken[0].scheduled_at.is_none());
4115
4116        let fetched = store.get_run(due.id).await.unwrap().unwrap();
4117        assert_eq!(fetched.status.state, RunStatus::Pending);
4118        assert!(fetched.scheduled_at.is_none());
4119
4120        // Exactly once: a second tick finds nothing.
4121        assert!(store.claim_due_sleeping_runs(10).await.unwrap().is_empty());
4122    }
4123
4124    #[tokio::test]
4125    async fn claim_due_sleeping_runs_skips_future_runs() {
4126        let store = InMemoryStore::new();
4127        let future = sleeping_run(&store, Utc::now() + TimeDelta::hours(1)).await;
4128
4129        assert!(store.claim_due_sleeping_runs(10).await.unwrap().is_empty());
4130        let fetched = store.get_run(future.id).await.unwrap().unwrap();
4131        assert_eq!(fetched.status.state, RunStatus::Sleeping);
4132    }
4133
4134    #[tokio::test]
4135    async fn claim_due_sleeping_runs_honours_limit_oldest_first() {
4136        let store = InMemoryStore::new();
4137        let older = sleeping_run(&store, Utc::now() - TimeDelta::seconds(20)).await;
4138        let newer = sleeping_run(&store, Utc::now() - TimeDelta::seconds(10)).await;
4139
4140        let woken = store.claim_due_sleeping_runs(1).await.unwrap();
4141        assert_eq!(woken.len(), 1);
4142        assert_eq!(woken[0].id, older.id);
4143
4144        let woken = store.claim_due_sleeping_runs(1).await.unwrap();
4145        assert_eq!(woken[0].id, newer.id);
4146    }
4147}