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