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