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