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