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