Skip to main content

ironflow_store/memory/
run_store.rs

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