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                },
2774                payload: json!({}),
2775                max_retries: 1,
2776                handler_version: None,
2777                labels: HashMap::new(),
2778                scheduled_at: None,
2779                idempotency_key: None,
2780                concurrency_key: None,
2781                priority: 0,
2782                concurrency_limits: Vec::new(),
2783                max_cost_usd: None,
2784                worker_tags: Vec::new(),
2785            })
2786            .await
2787            .unwrap()
2788            .into_run();
2789
2790        let r4 = store
2791            .create_run(NewRun {
2792                created_by: None,
2793                workflow_name: "test".to_string(),
2794                trigger: TriggerKind::Api,
2795                payload: json!({}),
2796                max_retries: 1,
2797                handler_version: None,
2798                labels: HashMap::new(),
2799                scheduled_at: None,
2800                idempotency_key: None,
2801                concurrency_key: None,
2802                priority: 0,
2803                concurrency_limits: Vec::new(),
2804                max_cost_usd: None,
2805                worker_tags: Vec::new(),
2806            })
2807            .await
2808            .unwrap()
2809            .into_run();
2810
2811        let r5 = store
2812            .create_run(NewRun {
2813                created_by: None,
2814                workflow_name: "test".to_string(),
2815                trigger: TriggerKind::Retry {
2816                    parent_run_id: Uuid::nil(),
2817                },
2818                payload: json!({}),
2819                max_retries: 1,
2820                handler_version: None,
2821                labels: HashMap::new(),
2822                scheduled_at: None,
2823                idempotency_key: None,
2824                concurrency_key: None,
2825                priority: 0,
2826                concurrency_limits: Vec::new(),
2827                max_cost_usd: None,
2828                worker_tags: Vec::new(),
2829            })
2830            .await
2831            .unwrap()
2832            .into_run();
2833
2834        assert_eq!(r1.trigger, TriggerKind::Manual);
2835        assert!(matches!(r2.trigger, TriggerKind::Webhook { .. }));
2836        assert!(matches!(r3.trigger, TriggerKind::Cron { .. }));
2837        assert_eq!(r4.trigger, TriggerKind::Api);
2838        assert!(matches!(r5.trigger, TriggerKind::Retry { .. }));
2839    }
2840
2841    // ---- create_step_dependencies ----
2842
2843    #[tokio::test]
2844    async fn create_step_dependencies_stores_dependencies() {
2845        let store = InMemoryStore::new();
2846        let run = store
2847            .create_run(new_run_req("test"))
2848            .await
2849            .unwrap()
2850            .into_run();
2851
2852        let step1 = store
2853            .create_step(NewStep {
2854                run_id: run.id,
2855                trace_id: step_trace_id(run.id, "step1", 0),
2856                name: "step1".to_string(),
2857                kind: crate::entities::StepKind::Shell,
2858                position: 0,
2859                input: None,
2860                is_error_handler: false,
2861            })
2862            .await
2863            .unwrap();
2864
2865        let step2 = store
2866            .create_step(NewStep {
2867                run_id: run.id,
2868                trace_id: step_trace_id(run.id, "step2", 1),
2869                name: "step2".to_string(),
2870                kind: crate::entities::StepKind::Shell,
2871                position: 1,
2872                input: None,
2873                is_error_handler: false,
2874            })
2875            .await
2876            .unwrap();
2877
2878        let result = store
2879            .create_step_dependencies(vec![NewStepDependency {
2880                step_id: step2.id,
2881                depends_on: step1.id,
2882            }])
2883            .await;
2884
2885        assert!(result.is_ok());
2886
2887        let deps = store.list_step_dependencies(run.id).await.unwrap();
2888        assert_eq!(deps.len(), 1);
2889        assert_eq!(deps[0].step_id, step2.id);
2890        assert_eq!(deps[0].depends_on, step1.id);
2891    }
2892
2893    #[tokio::test]
2894    async fn create_step_dependencies_duplicate_dependencies_are_idempotent() {
2895        let store = InMemoryStore::new();
2896        let run = store
2897            .create_run(new_run_req("test"))
2898            .await
2899            .unwrap()
2900            .into_run();
2901
2902        let step1 = store
2903            .create_step(NewStep {
2904                run_id: run.id,
2905                trace_id: step_trace_id(run.id, "step1", 0),
2906                name: "step1".to_string(),
2907                kind: crate::entities::StepKind::Shell,
2908                position: 0,
2909                input: None,
2910                is_error_handler: false,
2911            })
2912            .await
2913            .unwrap();
2914
2915        let step2 = store
2916            .create_step(NewStep {
2917                run_id: run.id,
2918                trace_id: step_trace_id(run.id, "step2", 1),
2919                name: "step2".to_string(),
2920                kind: crate::entities::StepKind::Shell,
2921                position: 1,
2922                input: None,
2923                is_error_handler: false,
2924            })
2925            .await
2926            .unwrap();
2927
2928        let dep = NewStepDependency {
2929            step_id: step2.id,
2930            depends_on: step1.id,
2931        };
2932
2933        store
2934            .create_step_dependencies(vec![dep.clone()])
2935            .await
2936            .unwrap();
2937        store.create_step_dependencies(vec![dep]).await.unwrap();
2938
2939        let deps = store.list_step_dependencies(run.id).await.unwrap();
2940        assert_eq!(deps.len(), 1);
2941    }
2942
2943    #[tokio::test]
2944    async fn create_step_dependencies_missing_step_id_returns_error() {
2945        let store = InMemoryStore::new();
2946        let run = store
2947            .create_run(new_run_req("test"))
2948            .await
2949            .unwrap()
2950            .into_run();
2951
2952        let step1 = store
2953            .create_step(NewStep {
2954                run_id: run.id,
2955                trace_id: step_trace_id(run.id, "step1", 0),
2956                name: "step1".to_string(),
2957                kind: crate::entities::StepKind::Shell,
2958                position: 0,
2959                input: None,
2960                is_error_handler: false,
2961            })
2962            .await
2963            .unwrap();
2964
2965        let result = store
2966            .create_step_dependencies(vec![NewStepDependency {
2967                step_id: Uuid::nil(),
2968                depends_on: step1.id,
2969            }])
2970            .await;
2971
2972        assert!(matches!(result.unwrap_err(), StoreError::StepNotFound(_)));
2973    }
2974
2975    #[tokio::test]
2976    async fn create_step_dependencies_missing_depends_on_returns_error() {
2977        let store = InMemoryStore::new();
2978        let run = store
2979            .create_run(new_run_req("test"))
2980            .await
2981            .unwrap()
2982            .into_run();
2983
2984        let step1 = store
2985            .create_step(NewStep {
2986                run_id: run.id,
2987                trace_id: step_trace_id(run.id, "step1", 0),
2988                name: "step1".to_string(),
2989                kind: crate::entities::StepKind::Shell,
2990                position: 0,
2991                input: None,
2992                is_error_handler: false,
2993            })
2994            .await
2995            .unwrap();
2996
2997        let result = store
2998            .create_step_dependencies(vec![NewStepDependency {
2999                step_id: step1.id,
3000                depends_on: Uuid::nil(),
3001            }])
3002            .await;
3003
3004        assert!(matches!(result.unwrap_err(), StoreError::StepNotFound(_)));
3005    }
3006
3007    #[tokio::test]
3008    async fn create_step_dependencies_multiple_dependencies() {
3009        let store = InMemoryStore::new();
3010        let run = store
3011            .create_run(new_run_req("test"))
3012            .await
3013            .unwrap()
3014            .into_run();
3015
3016        let step1 = store
3017            .create_step(NewStep {
3018                run_id: run.id,
3019                trace_id: step_trace_id(run.id, "step1", 0),
3020                name: "step1".to_string(),
3021                kind: crate::entities::StepKind::Shell,
3022                position: 0,
3023                input: None,
3024                is_error_handler: false,
3025            })
3026            .await
3027            .unwrap();
3028
3029        let step2 = store
3030            .create_step(NewStep {
3031                run_id: run.id,
3032                trace_id: step_trace_id(run.id, "step2", 1),
3033                name: "step2".to_string(),
3034                kind: crate::entities::StepKind::Shell,
3035                position: 1,
3036                input: None,
3037                is_error_handler: false,
3038            })
3039            .await
3040            .unwrap();
3041
3042        let step3 = store
3043            .create_step(NewStep {
3044                run_id: run.id,
3045                trace_id: step_trace_id(run.id, "step3", 2),
3046                name: "step3".to_string(),
3047                kind: crate::entities::StepKind::Shell,
3048                position: 2,
3049                input: None,
3050                is_error_handler: false,
3051            })
3052            .await
3053            .unwrap();
3054
3055        let result = store
3056            .create_step_dependencies(vec![
3057                NewStepDependency {
3058                    step_id: step2.id,
3059                    depends_on: step1.id,
3060                },
3061                NewStepDependency {
3062                    step_id: step3.id,
3063                    depends_on: step2.id,
3064                },
3065            ])
3066            .await;
3067
3068        assert!(result.is_ok());
3069
3070        let deps = store.list_step_dependencies(run.id).await.unwrap();
3071        assert_eq!(deps.len(), 2);
3072    }
3073
3074    // ---- list_step_dependencies ----
3075
3076    #[tokio::test]
3077    async fn list_step_dependencies_empty_for_run_with_no_dependencies() {
3078        let store = InMemoryStore::new();
3079        let run = store
3080            .create_run(new_run_req("test"))
3081            .await
3082            .unwrap()
3083            .into_run();
3084
3085        store
3086            .create_step(NewStep {
3087                run_id: run.id,
3088                trace_id: step_trace_id(run.id, "step1", 0),
3089                name: "step1".to_string(),
3090                kind: crate::entities::StepKind::Shell,
3091                position: 0,
3092                input: None,
3093                is_error_handler: false,
3094            })
3095            .await
3096            .unwrap();
3097
3098        let deps = store.list_step_dependencies(run.id).await.unwrap();
3099        assert!(deps.is_empty());
3100    }
3101
3102    #[tokio::test]
3103    async fn list_step_dependencies_returns_only_deps_for_given_run() {
3104        let store = InMemoryStore::new();
3105        let run1 = store
3106            .create_run(new_run_req("test1"))
3107            .await
3108            .unwrap()
3109            .into_run();
3110        let run2 = store
3111            .create_run(new_run_req("test2"))
3112            .await
3113            .unwrap()
3114            .into_run();
3115
3116        let step1_run1 = store
3117            .create_step(NewStep {
3118                run_id: run1.id,
3119                trace_id: step_trace_id(run1.id, "step1", 0),
3120                name: "step1".to_string(),
3121                kind: crate::entities::StepKind::Shell,
3122                position: 0,
3123                input: None,
3124                is_error_handler: false,
3125            })
3126            .await
3127            .unwrap();
3128
3129        let step2_run1 = store
3130            .create_step(NewStep {
3131                run_id: run1.id,
3132                trace_id: step_trace_id(run1.id, "step2", 1),
3133                name: "step2".to_string(),
3134                kind: crate::entities::StepKind::Shell,
3135                position: 1,
3136                input: None,
3137                is_error_handler: false,
3138            })
3139            .await
3140            .unwrap();
3141
3142        let step1_run2 = store
3143            .create_step(NewStep {
3144                run_id: run2.id,
3145                trace_id: step_trace_id(run2.id, "step1", 0),
3146                name: "step1".to_string(),
3147                kind: crate::entities::StepKind::Shell,
3148                position: 0,
3149                input: None,
3150                is_error_handler: false,
3151            })
3152            .await
3153            .unwrap();
3154
3155        let step2_run2 = store
3156            .create_step(NewStep {
3157                run_id: run2.id,
3158                trace_id: step_trace_id(run2.id, "step2", 1),
3159                name: "step2".to_string(),
3160                kind: crate::entities::StepKind::Shell,
3161                position: 1,
3162                input: None,
3163                is_error_handler: false,
3164            })
3165            .await
3166            .unwrap();
3167
3168        store
3169            .create_step_dependencies(vec![
3170                NewStepDependency {
3171                    step_id: step2_run1.id,
3172                    depends_on: step1_run1.id,
3173                },
3174                NewStepDependency {
3175                    step_id: step2_run2.id,
3176                    depends_on: step1_run2.id,
3177                },
3178            ])
3179            .await
3180            .unwrap();
3181
3182        let deps_run1 = store.list_step_dependencies(run1.id).await.unwrap();
3183        let deps_run2 = store.list_step_dependencies(run2.id).await.unwrap();
3184
3185        assert_eq!(deps_run1.len(), 1);
3186        assert_eq!(deps_run1[0].step_id, step2_run1.id);
3187        assert_eq!(deps_run1[0].depends_on, step1_run1.id);
3188
3189        assert_eq!(deps_run2.len(), 1);
3190        assert_eq!(deps_run2[0].step_id, step2_run2.id);
3191        assert_eq!(deps_run2[0].depends_on, step1_run2.id);
3192    }
3193
3194    #[tokio::test]
3195    async fn list_step_dependencies_returns_empty_for_nonexistent_run() {
3196        let store = InMemoryStore::new();
3197        let deps = store.list_step_dependencies(Uuid::nil()).await.unwrap();
3198        assert!(deps.is_empty());
3199    }
3200
3201    #[tokio::test]
3202    async fn list_step_dependencies_sorted_by_created_at() {
3203        let store = InMemoryStore::new();
3204        let run = store
3205            .create_run(new_run_req("test"))
3206            .await
3207            .unwrap()
3208            .into_run();
3209
3210        let step1 = store
3211            .create_step(NewStep {
3212                run_id: run.id,
3213                trace_id: step_trace_id(run.id, "step1", 0),
3214                name: "step1".to_string(),
3215                kind: crate::entities::StepKind::Shell,
3216                position: 0,
3217                input: None,
3218                is_error_handler: false,
3219            })
3220            .await
3221            .unwrap();
3222
3223        let step2 = store
3224            .create_step(NewStep {
3225                run_id: run.id,
3226                trace_id: step_trace_id(run.id, "step2", 1),
3227                name: "step2".to_string(),
3228                kind: crate::entities::StepKind::Shell,
3229                position: 1,
3230                input: None,
3231                is_error_handler: false,
3232            })
3233            .await
3234            .unwrap();
3235
3236        let step3 = store
3237            .create_step(NewStep {
3238                run_id: run.id,
3239                trace_id: step_trace_id(run.id, "step3", 2),
3240                name: "step3".to_string(),
3241                kind: crate::entities::StepKind::Shell,
3242                position: 2,
3243                input: None,
3244                is_error_handler: false,
3245            })
3246            .await
3247            .unwrap();
3248
3249        store
3250            .create_step_dependencies(vec![NewStepDependency {
3251                step_id: step2.id,
3252                depends_on: step1.id,
3253            }])
3254            .await
3255            .unwrap();
3256
3257        store
3258            .create_step_dependencies(vec![NewStepDependency {
3259                step_id: step3.id,
3260                depends_on: step1.id,
3261            }])
3262            .await
3263            .unwrap();
3264
3265        let deps = store.list_step_dependencies(run.id).await.unwrap();
3266        assert_eq!(deps.len(), 2);
3267        assert!(deps[0].created_at <= deps[1].created_at);
3268    }
3269
3270    // ---- update_run_returning ----
3271
3272    #[tokio::test]
3273    async fn update_run_returning_applies_and_returns() {
3274        let store = InMemoryStore::new();
3275        let run = store
3276            .create_run(new_run_req("test"))
3277            .await
3278            .unwrap()
3279            .into_run();
3280
3281        // Transition Pending -> Running first
3282        store
3283            .update_run_status(run.id, RunStatus::Running)
3284            .await
3285            .unwrap();
3286
3287        let updated = store
3288            .update_run_returning(
3289                run.id,
3290                RunUpdate {
3291                    status: Some(RunStatus::Completed),
3292                    cost_usd: Some(Decimal::new(4200, 2)),
3293                    duration_ms: Some(1500),
3294                    ..RunUpdate::default()
3295                },
3296            )
3297            .await
3298            .unwrap();
3299
3300        assert_eq!(updated.id, run.id);
3301        assert_eq!(updated.status.state, RunStatus::Completed);
3302        assert_eq!(updated.cost_usd, Decimal::new(4200, 2));
3303        assert_eq!(updated.duration_ms, 1500);
3304        assert!(updated.completed_at.is_some());
3305    }
3306
3307    #[tokio::test]
3308    async fn update_run_returning_not_found() {
3309        let store = InMemoryStore::new();
3310        let result = store
3311            .update_run_returning(
3312                Uuid::nil(),
3313                RunUpdate {
3314                    status: Some(RunStatus::Running),
3315                    ..RunUpdate::default()
3316                },
3317            )
3318            .await;
3319
3320        assert!(matches!(result, Err(StoreError::RunNotFound(_))));
3321    }
3322
3323    #[tokio::test]
3324    async fn update_run_returning_invalid_transition() {
3325        let store = InMemoryStore::new();
3326        let run = store
3327            .create_run(new_run_req("test"))
3328            .await
3329            .unwrap()
3330            .into_run();
3331
3332        let result = store
3333            .update_run_returning(
3334                run.id,
3335                RunUpdate {
3336                    status: Some(RunStatus::Completed),
3337                    ..RunUpdate::default()
3338                },
3339            )
3340            .await;
3341
3342        assert!(matches!(result, Err(StoreError::InvalidTransition { .. })));
3343    }
3344
3345    // ---- retry scheduling ----
3346
3347    #[tokio::test]
3348    async fn create_step_stamps_the_current_attempt() {
3349        let store = InMemoryStore::new();
3350        let run = store
3351            .create_run(new_run_req("retry-wf"))
3352            .await
3353            .unwrap()
3354            .into_run();
3355
3356        let first = store
3357            .create_step(new_step_req(run.id, "build", 0))
3358            .await
3359            .unwrap();
3360        assert_eq!(first.attempt, 1);
3361
3362        store
3363            .update_run_status(run.id, RunStatus::Running)
3364            .await
3365            .unwrap();
3366        store
3367            .update_run(
3368                run.id,
3369                RunUpdate {
3370                    status: Some(RunStatus::Retrying),
3371                    increment_retry: true,
3372                    ..RunUpdate::default()
3373                },
3374            )
3375            .await
3376            .unwrap();
3377
3378        let second = store
3379            .create_step(new_step_req(run.id, "build", 0))
3380            .await
3381            .unwrap();
3382        assert_eq!(second.attempt, 2);
3383    }
3384
3385    #[tokio::test]
3386    async fn pick_next_pending_ignores_retrying_run_before_its_backoff() {
3387        let store = InMemoryStore::new();
3388        let run = store
3389            .create_run(new_run_req("retry-wf"))
3390            .await
3391            .unwrap()
3392            .into_run();
3393
3394        store
3395            .update_run_status(run.id, RunStatus::Running)
3396            .await
3397            .unwrap();
3398        store
3399            .update_run(
3400                run.id,
3401                RunUpdate {
3402                    status: Some(RunStatus::Retrying),
3403                    increment_retry: true,
3404                    scheduled_at: Some(Utc::now() + TimeDelta::seconds(60)),
3405                    ..RunUpdate::default()
3406                },
3407            )
3408            .await
3409            .unwrap();
3410
3411        assert!(store.pick_next_pending(None).await.unwrap().is_none());
3412    }
3413
3414    #[tokio::test]
3415    async fn pick_next_pending_resumes_retrying_run_after_its_backoff() {
3416        let store = InMemoryStore::new();
3417        let run = store
3418            .create_run(new_run_req("retry-wf"))
3419            .await
3420            .unwrap()
3421            .into_run();
3422
3423        store
3424            .update_run_status(run.id, RunStatus::Running)
3425            .await
3426            .unwrap();
3427        store
3428            .update_run(
3429                run.id,
3430                RunUpdate {
3431                    status: Some(RunStatus::Retrying),
3432                    increment_retry: true,
3433                    scheduled_at: Some(Utc::now() - TimeDelta::seconds(1)),
3434                    ..RunUpdate::default()
3435                },
3436            )
3437            .await
3438            .unwrap();
3439
3440        let picked = store.pick_next_pending(None).await.unwrap().unwrap();
3441        assert_eq!(picked.id, run.id);
3442        assert_eq!(picked.status.state, RunStatus::Running);
3443        assert_eq!(picked.retry_count, 1);
3444    }
3445
3446    // ---- priority ----
3447
3448    async fn create_with_priority(store: &InMemoryStore, name: &str, priority: i16) -> Run {
3449        let run = store
3450            .create_run(NewRun {
3451                priority,
3452                ..new_run_req(name)
3453            })
3454            .await
3455            .unwrap()
3456            .into_run();
3457        // Distinct created_at values, so FIFO among equal priorities is observable.
3458        sleep(Duration::from_millis(2)).await;
3459        run
3460    }
3461
3462    #[tokio::test]
3463    async fn pick_next_pending_priority_serves_higher_priority_first() {
3464        let store = InMemoryStore::new();
3465        let low = create_with_priority(&store, "low", 0).await;
3466        let high = create_with_priority(&store, "high", 10).await;
3467
3468        let first = store.pick_next_pending(None).await.unwrap().unwrap();
3469        assert_eq!(first.id, high.id);
3470        assert_eq!(first.priority, 10);
3471        let second = store.pick_next_pending(None).await.unwrap().unwrap();
3472        assert_eq!(second.id, low.id);
3473    }
3474
3475    #[tokio::test]
3476    async fn pick_next_pending_priority_is_fifo_among_equal_priorities() {
3477        let store = InMemoryStore::new();
3478        let older = create_with_priority(&store, "older", 5).await;
3479        let younger = create_with_priority(&store, "younger", 5).await;
3480
3481        let first = store.pick_next_pending(None).await.unwrap().unwrap();
3482        assert_eq!(first.id, older.id);
3483        let second = store.pick_next_pending(None).await.unwrap().unwrap();
3484        assert_eq!(second.id, younger.id);
3485    }
3486
3487    #[tokio::test]
3488    async fn pick_next_pending_priority_negative_runs_after_default() {
3489        let store = InMemoryStore::new();
3490        let negative = create_with_priority(&store, "background", -50).await;
3491        let default = create_with_priority(&store, "default", 0).await;
3492
3493        let first = store.pick_next_pending(None).await.unwrap().unwrap();
3494        assert_eq!(first.id, default.id);
3495        let second = store.pick_next_pending(None).await.unwrap().unwrap();
3496        assert_eq!(second.id, negative.id);
3497    }
3498
3499    #[tokio::test]
3500    async fn pick_next_pending_priority_skips_high_priority_run_not_yet_due() {
3501        let store = InMemoryStore::new();
3502        let later = store
3503            .create_run(NewRun {
3504                priority: 100,
3505                scheduled_at: Some(Utc::now() + TimeDelta::seconds(3600)),
3506                ..new_run_req("later")
3507            })
3508            .await
3509            .unwrap()
3510            .into_run();
3511        let now = create_with_priority(&store, "now", 0).await;
3512
3513        let picked = store.pick_next_pending(None).await.unwrap().unwrap();
3514        assert_eq!(picked.id, now.id);
3515        assert!(store.pick_next_pending(None).await.unwrap().is_none());
3516        let later = store.get_run(later.id).await.unwrap().unwrap();
3517        assert_eq!(later.status.state, RunStatus::Pending);
3518    }
3519
3520    #[tokio::test]
3521    async fn pick_next_pending_priority_kept_after_retry() {
3522        let store = InMemoryStore::new();
3523        let run = create_with_priority(&store, "retry-wf", 42).await;
3524
3525        store
3526            .update_run_status(run.id, RunStatus::Running)
3527            .await
3528            .unwrap();
3529        store
3530            .update_run(
3531                run.id,
3532                RunUpdate {
3533                    status: Some(RunStatus::Retrying),
3534                    increment_retry: true,
3535                    scheduled_at: Some(Utc::now() - TimeDelta::seconds(1)),
3536                    ..RunUpdate::default()
3537                },
3538            )
3539            .await
3540            .unwrap();
3541        let fresh = create_with_priority(&store, "fresh", 0).await;
3542
3543        let picked = store.pick_next_pending(None).await.unwrap().unwrap();
3544        assert_eq!(picked.id, run.id);
3545        assert_eq!(picked.priority, 42);
3546        assert_eq!(picked.retry_count, 1);
3547        let next = store.pick_next_pending(None).await.unwrap().unwrap();
3548        assert_eq!(next.id, fresh.id);
3549    }
3550
3551    #[tokio::test]
3552    async fn create_run_priority_out_of_range_is_rejected() {
3553        let store = InMemoryStore::new();
3554        for priority in [101, -101] {
3555            let err = store
3556                .create_run(NewRun {
3557                    priority,
3558                    ..new_run_req("out-of-range")
3559                })
3560                .await
3561                .unwrap_err();
3562            assert!(matches!(err, StoreError::Database(_)), "{err:?}");
3563        }
3564        let page = store.list_runs(RunFilter::default(), 1, 10).await.unwrap();
3565        assert_eq!(page.total, 0);
3566    }
3567
3568    #[tokio::test]
3569    async fn list_runs_priority_filter_is_exact_match() {
3570        let store = InMemoryStore::new();
3571        let urgent = create_with_priority(&store, "urgent", 10).await;
3572        create_with_priority(&store, "default", 0).await;
3573        create_with_priority(&store, "more-urgent", 20).await;
3574
3575        let page = store
3576            .list_runs(
3577                RunFilter {
3578                    priority: Some(10),
3579                    ..RunFilter::default()
3580                },
3581                1,
3582                10,
3583            )
3584            .await
3585            .unwrap();
3586        assert_eq!(page.total, 1);
3587        assert_eq!(page.items[0].id, urgent.id);
3588
3589        let all = store.list_runs(RunFilter::default(), 1, 10).await.unwrap();
3590        assert_eq!(all.total, 3);
3591    }
3592
3593    #[tokio::test]
3594    async fn update_run_persists_scheduled_at() {
3595        let store = InMemoryStore::new();
3596        let run = store
3597            .create_run(new_run_req("test"))
3598            .await
3599            .unwrap()
3600            .into_run();
3601        let when = Utc::now() + TimeDelta::seconds(30);
3602
3603        store
3604            .update_run(
3605                run.id,
3606                RunUpdate {
3607                    scheduled_at: Some(when),
3608                    ..RunUpdate::default()
3609                },
3610            )
3611            .await
3612            .unwrap();
3613
3614        let fetched = store.get_run(run.id).await.unwrap().unwrap();
3615        assert_eq!(fetched.scheduled_at, Some(when));
3616    }
3617
3618    // ---- output ----
3619
3620    #[tokio::test]
3621    async fn new_run_has_no_output() {
3622        let store = InMemoryStore::new();
3623        let run = store
3624            .create_run(new_run_req("test"))
3625            .await
3626            .unwrap()
3627            .into_run();
3628
3629        assert!(run.output.is_none());
3630        let fetched = store.get_run(run.id).await.unwrap().unwrap();
3631        assert!(fetched.output.is_none());
3632    }
3633
3634    #[tokio::test]
3635    async fn update_run_sets_output() {
3636        let store = InMemoryStore::new();
3637        let run = store
3638            .create_run(new_run_req("test"))
3639            .await
3640            .unwrap()
3641            .into_run();
3642
3643        store
3644            .update_run(
3645                run.id,
3646                RunUpdate {
3647                    output: Some(json!({"verdict": "approved"})),
3648                    ..RunUpdate::default()
3649                },
3650            )
3651            .await
3652            .unwrap();
3653
3654        let fetched = store.get_run(run.id).await.unwrap().unwrap();
3655        assert_eq!(fetched.output, Some(json!({"verdict": "approved"})));
3656    }
3657
3658    #[tokio::test]
3659    async fn update_run_without_output_keeps_previous_output() {
3660        let store = InMemoryStore::new();
3661        let run = store
3662            .create_run(new_run_req("test"))
3663            .await
3664            .unwrap()
3665            .into_run();
3666
3667        store
3668            .update_run(
3669                run.id,
3670                RunUpdate {
3671                    output: Some(json!({"verdict": "approved"})),
3672                    ..RunUpdate::default()
3673                },
3674            )
3675            .await
3676            .unwrap();
3677        store
3678            .update_run(
3679                run.id,
3680                RunUpdate {
3681                    error: Some("boom".to_string()),
3682                    ..RunUpdate::default()
3683                },
3684            )
3685            .await
3686            .unwrap();
3687
3688        let fetched = store.get_run(run.id).await.unwrap().unwrap();
3689        assert_eq!(fetched.output, Some(json!({"verdict": "approved"})));
3690        assert_eq!(fetched.error.as_deref(), Some("boom"));
3691    }
3692
3693    #[tokio::test]
3694    async fn update_run_output_last_write_wins() {
3695        let store = InMemoryStore::new();
3696        let run = store
3697            .create_run(new_run_req("test"))
3698            .await
3699            .unwrap()
3700            .into_run();
3701
3702        for verdict in ["first", "second"] {
3703            store
3704                .update_run(
3705                    run.id,
3706                    RunUpdate {
3707                        output: Some(json!({ "verdict": verdict })),
3708                        ..RunUpdate::default()
3709                    },
3710                )
3711                .await
3712                .unwrap();
3713        }
3714
3715        let fetched = store.get_run(run.id).await.unwrap().unwrap();
3716        assert_eq!(fetched.output, Some(json!({"verdict": "second"})));
3717    }
3718
3719    // ---- created_by ----
3720
3721    async fn seed_user(store: &InMemoryStore, username: &str) -> Uuid {
3722        store
3723            .create_user(NewUser {
3724                email: format!("{username}@example.com"),
3725                username: username.to_string(),
3726                password_hash: "hash".to_string(),
3727                is_admin: Some(false),
3728            })
3729            .await
3730            .unwrap()
3731            .id
3732    }
3733
3734    async fn seed_api_key(store: &InMemoryStore, user_id: Uuid, name: &str) -> Uuid {
3735        store
3736            .create_api_key(NewApiKey {
3737                user_id,
3738                name: name.to_string(),
3739                key_hash: "hash".to_string(),
3740                key_prefix: "irfl_0000".to_string(),
3741                scopes: vec![ApiKeyScope::RunsWrite],
3742                expires_at: None,
3743                rate_limit_override: None,
3744            })
3745            .await
3746            .unwrap()
3747            .id
3748    }
3749
3750    fn run_req_by(actor: RunActor) -> NewRun {
3751        NewRun {
3752            created_by: Some(actor),
3753            ..new_run_req("test")
3754        }
3755    }
3756
3757    #[tokio::test]
3758    async fn create_run_without_actor_has_no_author() {
3759        let store = InMemoryStore::new();
3760        let run = store
3761            .create_run(new_run_req("test"))
3762            .await
3763            .unwrap()
3764            .into_run();
3765
3766        assert!(run.created_by.is_none());
3767        assert!(run.created_by_label.is_none());
3768    }
3769
3770    #[tokio::test]
3771    async fn create_run_by_user_resolves_username_as_label() {
3772        let store = InMemoryStore::new();
3773        let user_id = seed_user(&store, "alice").await;
3774
3775        let run = store
3776            .create_run(run_req_by(RunActor::User { user_id }))
3777            .await
3778            .unwrap()
3779            .into_run();
3780
3781        assert_eq!(run.created_by, Some(RunActor::User { user_id }));
3782        assert_eq!(run.created_by_label.as_deref(), Some("alice"));
3783    }
3784
3785    #[tokio::test]
3786    async fn create_run_by_api_key_resolves_key_and_owner_as_label() {
3787        let store = InMemoryStore::new();
3788        let user_id = seed_user(&store, "alice").await;
3789        let api_key_id = seed_api_key(&store, user_id, "ci-deploy").await;
3790
3791        let run = store
3792            .create_run(run_req_by(RunActor::ApiKey {
3793                api_key_id,
3794                user_id,
3795            }))
3796            .await
3797            .unwrap()
3798            .into_run();
3799
3800        assert_eq!(run.created_by_label.as_deref(), Some("ci-deploy (alice)"));
3801    }
3802
3803    #[tokio::test]
3804    async fn label_follows_api_key_rename() {
3805        let store = InMemoryStore::new();
3806        let user_id = seed_user(&store, "alice").await;
3807        let api_key_id = seed_api_key(&store, user_id, "ci-deploy").await;
3808        let run = store
3809            .create_run(run_req_by(RunActor::ApiKey {
3810                api_key_id,
3811                user_id,
3812            }))
3813            .await
3814            .unwrap()
3815            .into_run();
3816
3817        store
3818            .update_api_key(
3819                api_key_id,
3820                ApiKeyUpdate {
3821                    name: Some("ci-release".to_string()),
3822                    ..ApiKeyUpdate::default()
3823                },
3824            )
3825            .await
3826            .unwrap();
3827
3828        let reread = store.get_run(run.id).await.unwrap().unwrap();
3829        assert_eq!(
3830            reread.created_by_label.as_deref(),
3831            Some("ci-release (alice)")
3832        );
3833    }
3834
3835    #[tokio::test]
3836    async fn label_is_none_when_user_is_unknown() {
3837        let store = InMemoryStore::new();
3838        let run = store
3839            .create_run(run_req_by(RunActor::User {
3840                user_id: Uuid::now_v7(),
3841            }))
3842            .await
3843            .unwrap()
3844            .into_run();
3845
3846        assert!(run.created_by.is_some());
3847        assert!(run.created_by_label.is_none());
3848    }
3849
3850    #[tokio::test]
3851    async fn label_is_key_name_only_when_owner_is_unknown() {
3852        let store = InMemoryStore::new();
3853        let owner = seed_user(&store, "alice").await;
3854        let api_key_id = seed_api_key(&store, owner, "ci-deploy").await;
3855
3856        // An actor pointing at a user that was never created.
3857        let run = store
3858            .create_run(run_req_by(RunActor::ApiKey {
3859                api_key_id,
3860                user_id: Uuid::now_v7(),
3861            }))
3862            .await
3863            .unwrap()
3864            .into_run();
3865
3866        assert_eq!(run.created_by_label.as_deref(), Some("ci-deploy"));
3867    }
3868
3869    #[tokio::test]
3870    async fn list_runs_filters_by_author() {
3871        let store = InMemoryStore::new();
3872        let alice = seed_user(&store, "alice").await;
3873        let bob = seed_user(&store, "bob").await;
3874
3875        store
3876            .create_run(run_req_by(RunActor::User { user_id: alice }))
3877            .await
3878            .unwrap()
3879            .into_run();
3880        store
3881            .create_run(run_req_by(RunActor::User { user_id: bob }))
3882            .await
3883            .unwrap()
3884            .into_run();
3885        store.create_run(new_run_req("anonymous")).await.unwrap();
3886
3887        let page = store
3888            .list_runs(
3889                RunFilter {
3890                    created_by_user_id: Some(alice),
3891                    ..RunFilter::default()
3892                },
3893                1,
3894                20,
3895            )
3896            .await
3897            .unwrap();
3898
3899        assert_eq!(page.total, 1);
3900        assert_eq!(page.items[0].created_by_label.as_deref(), Some("alice"));
3901    }
3902
3903    #[tokio::test]
3904    async fn list_runs_author_filter_matches_runs_from_the_users_api_keys() {
3905        let store = InMemoryStore::new();
3906        let alice = seed_user(&store, "alice").await;
3907        let api_key_id = seed_api_key(&store, alice, "ci-deploy").await;
3908
3909        store
3910            .create_run(run_req_by(RunActor::ApiKey {
3911                api_key_id,
3912                user_id: alice,
3913            }))
3914            .await
3915            .unwrap()
3916            .into_run();
3917
3918        let page = store
3919            .list_runs(
3920                RunFilter {
3921                    created_by_user_id: Some(alice),
3922                    ..RunFilter::default()
3923                },
3924                1,
3925                20,
3926            )
3927            .await
3928            .unwrap();
3929
3930        assert_eq!(page.total, 1);
3931    }
3932
3933    #[tokio::test]
3934    async fn list_runs_author_filter_excludes_unrelated_users() {
3935        let store = InMemoryStore::new();
3936        let alice = seed_user(&store, "alice").await;
3937
3938        store
3939            .create_run(run_req_by(RunActor::User { user_id: alice }))
3940            .await
3941            .unwrap()
3942            .into_run();
3943
3944        let page = store
3945            .list_runs(
3946                RunFilter {
3947                    created_by_user_id: Some(Uuid::now_v7()),
3948                    ..RunFilter::default()
3949                },
3950                1,
3951                20,
3952            )
3953            .await
3954            .unwrap();
3955
3956        assert_eq!(page.total, 0);
3957    }
3958
3959    #[tokio::test]
3960    async fn list_runs_without_author_filter_returns_every_run() {
3961        let store = InMemoryStore::new();
3962        let alice = seed_user(&store, "alice").await;
3963
3964        store
3965            .create_run(run_req_by(RunActor::User { user_id: alice }))
3966            .await
3967            .unwrap()
3968            .into_run();
3969        store.create_run(new_run_req("anonymous")).await.unwrap();
3970
3971        let page = store.list_runs(RunFilter::default(), 1, 20).await.unwrap();
3972        assert_eq!(page.total, 2);
3973    }
3974
3975    #[tokio::test]
3976    async fn pick_next_pending_resolves_author_label() {
3977        let store = InMemoryStore::new();
3978        let user_id = seed_user(&store, "alice").await;
3979        store
3980            .create_run(run_req_by(RunActor::User { user_id }))
3981            .await
3982            .unwrap()
3983            .into_run();
3984
3985        let picked = store.pick_next_pending(None).await.unwrap().unwrap();
3986        assert_eq!(picked.created_by_label.as_deref(), Some("alice"));
3987    }
3988
3989    // ---- list_purgeable_runs ----
3990
3991    #[tokio::test]
3992    async fn list_purgeable_runs_returns_old_terminal_runs() {
3993        let store = InMemoryStore::new();
3994        let old = create_terminal_run(&store, "deploy", RunStatus::Completed).await;
3995        store
3996            .set_run_created_at(old.id, Utc::now() - chrono::Duration::days(100))
3997            .await;
3998
3999        let policy = PurgePolicy {
4000            max_age_days: 90,
4001            max_runs_per_workflow: 10000,
4002            dry_run: false,
4003        };
4004        let result = store.list_purgeable_runs(&policy, 100).await.unwrap();
4005
4006        assert_eq!(result.len(), 1);
4007        assert_eq!(result[0].run_id, old.id);
4008        assert_eq!(result[0].reason, PurgeReason::TooOld);
4009    }
4010
4011    #[tokio::test]
4012    async fn list_purgeable_runs_ignores_non_terminal_states() {
4013        let store = InMemoryStore::new();
4014
4015        // Pending
4016        let pending = store
4017            .create_run(new_run_req("deploy"))
4018            .await
4019            .unwrap()
4020            .into_run();
4021        store
4022            .set_run_created_at(pending.id, Utc::now() - chrono::Duration::days(200))
4023            .await;
4024
4025        // Running
4026        let running = store
4027            .create_run(new_run_req("deploy"))
4028            .await
4029            .unwrap()
4030            .into_run();
4031        store
4032            .update_run_status(running.id, RunStatus::Running)
4033            .await
4034            .unwrap();
4035        store
4036            .set_run_created_at(running.id, Utc::now() - chrono::Duration::days(200))
4037            .await;
4038
4039        let policy = PurgePolicy {
4040            max_age_days: 90,
4041            max_runs_per_workflow: 1,
4042            dry_run: false,
4043        };
4044        let result = store.list_purgeable_runs(&policy, 100).await.unwrap();
4045        assert!(result.is_empty());
4046    }
4047
4048    #[tokio::test]
4049    async fn list_purgeable_runs_returns_excess_per_workflow() {
4050        let store = InMemoryStore::new();
4051        let r1 = create_terminal_run(&store, "deploy", RunStatus::Completed).await;
4052        store
4053            .set_run_created_at(r1.id, Utc::now() - chrono::Duration::days(10))
4054            .await;
4055        let r2 = create_terminal_run(&store, "deploy", RunStatus::Completed).await;
4056        store
4057            .set_run_created_at(r2.id, Utc::now() - chrono::Duration::days(5))
4058            .await;
4059        let _r3 = create_terminal_run(&store, "deploy", RunStatus::Completed).await;
4060
4061        let policy = PurgePolicy {
4062            max_age_days: 365,
4063            max_runs_per_workflow: 2,
4064            dry_run: false,
4065        };
4066        let result = store.list_purgeable_runs(&policy, 100).await.unwrap();
4067
4068        assert_eq!(result.len(), 1);
4069        assert_eq!(result[0].run_id, r1.id);
4070        assert_eq!(result[0].reason, PurgeReason::ExceedsWorkflowLimit);
4071    }
4072
4073    // ---- delete_run ----
4074
4075    #[tokio::test]
4076    async fn delete_run_removes_run_and_associated_data() {
4077        use crate::artifact_store::ArtifactStore;
4078        use crate::entities::{NewStep, StepKind};
4079
4080        let store = InMemoryStore::new();
4081        let run = create_terminal_run(&store, "deploy", RunStatus::Completed).await;
4082        let step = store
4083            .create_step(NewStep {
4084                run_id: run.id,
4085                trace_id: step_trace_id(run.id, "build", 0),
4086                name: "build".to_string(),
4087                kind: StepKind::Shell,
4088                position: 0,
4089                input: None,
4090                is_error_handler: false,
4091            })
4092            .await
4093            .unwrap();
4094
4095        let artifact_id = Uuid::now_v7();
4096        store
4097            .create_artifact(crate::entities::NewArtifact {
4098                id: artifact_id,
4099                run_id: run.id,
4100                step_id: step.id,
4101                name: "report.html".to_string(),
4102                storage_key: format!("artifacts/{}/{}/{}", run.id, step.id, artifact_id),
4103                content_type: "text/html".to_string(),
4104                size_bytes: 42,
4105                sha256: "0".repeat(64),
4106            })
4107            .await
4108            .unwrap();
4109
4110        let keys = store.delete_run(run.id).await.unwrap();
4111
4112        assert_eq!(keys.len(), 1);
4113        assert!(keys[0].contains(&artifact_id.to_string()));
4114        assert!(store.get_run(run.id).await.unwrap().is_none());
4115        assert!(store.list_steps(run.id).await.unwrap().is_empty());
4116        assert!(
4117            store
4118                .list_artifacts_for_run(run.id)
4119                .await
4120                .unwrap()
4121                .is_empty()
4122        );
4123    }
4124
4125    #[tokio::test]
4126    async fn delete_run_not_found() {
4127        let store = InMemoryStore::new();
4128        let err = store.delete_run(Uuid::now_v7()).await.unwrap_err();
4129        assert!(matches!(err, StoreError::RunNotFound(_)));
4130    }
4131
4132    // ---- record_step_approval ----
4133
4134    fn vote(user_id: Uuid, name: &str) -> StepApproval {
4135        StepApproval {
4136            user_id,
4137            approved_by: name.to_string(),
4138            at: Utc::now(),
4139        }
4140    }
4141
4142    #[tokio::test]
4143    async fn record_step_approval_appends_distinct_voters() {
4144        let store = InMemoryStore::new();
4145        let run = store
4146            .create_run(new_run_req("test"))
4147            .await
4148            .unwrap()
4149            .into_run();
4150        let step = store
4151            .create_step(new_step_req(run.id, "gate", 0))
4152            .await
4153            .unwrap();
4154        assert!(step.approvals.is_empty());
4155        assert!(step.approval_requirement.is_none());
4156
4157        let alice = Uuid::now_v7();
4158        let bob = Uuid::now_v7();
4159        let after_first = store
4160            .record_step_approval(step.id, vote(alice, "alice"))
4161            .await
4162            .unwrap();
4163        assert_eq!(after_first.approvals.len(), 1);
4164
4165        let after_second = store
4166            .record_step_approval(step.id, vote(bob, "bob"))
4167            .await
4168            .unwrap();
4169        assert_eq!(after_second.approvals.len(), 2);
4170        assert_eq!(after_second.approvals[0].user_id, alice);
4171        assert_eq!(after_second.approvals[1].user_id, bob);
4172    }
4173
4174    #[tokio::test]
4175    async fn record_step_approval_ignores_same_user() {
4176        let store = InMemoryStore::new();
4177        let run = store
4178            .create_run(new_run_req("test"))
4179            .await
4180            .unwrap()
4181            .into_run();
4182        let step = store
4183            .create_step(new_step_req(run.id, "gate", 0))
4184            .await
4185            .unwrap();
4186
4187        let alice = Uuid::now_v7();
4188        store
4189            .record_step_approval(step.id, vote(alice, "alice"))
4190            .await
4191            .unwrap();
4192        let again = store
4193            .record_step_approval(step.id, vote(alice, "alice-key"))
4194            .await
4195            .unwrap();
4196
4197        assert_eq!(again.approvals.len(), 1);
4198        assert_eq!(again.approvals[0].approved_by, "alice");
4199    }
4200
4201    #[tokio::test]
4202    async fn record_step_approval_unknown_step_is_not_found() {
4203        let store = InMemoryStore::new();
4204        let err = store
4205            .record_step_approval(Uuid::now_v7(), vote(Uuid::now_v7(), "alice"))
4206            .await
4207            .unwrap_err();
4208        assert!(matches!(err, StoreError::StepNotFound(_)));
4209    }
4210
4211    #[tokio::test]
4212    async fn update_step_sets_approval_requirement() {
4213        let store = InMemoryStore::new();
4214        let run = store
4215            .create_run(new_run_req("test"))
4216            .await
4217            .unwrap()
4218            .into_run();
4219        let step = store
4220            .create_step(new_step_req(run.id, "gate", 0))
4221            .await
4222            .unwrap();
4223        let requirement = ApprovalRequirement {
4224            required_approvers: 3,
4225            ..ApprovalRequirement::default()
4226        };
4227
4228        store
4229            .update_step(
4230                step.id,
4231                StepUpdate {
4232                    approval_requirement: Some(requirement.clone()),
4233                    ..StepUpdate::default()
4234                },
4235            )
4236            .await
4237            .unwrap();
4238
4239        let fetched = store.get_step(step.id).await.unwrap().unwrap();
4240        assert_eq!(fetched.approval_requirement, Some(requirement));
4241    }
4242
4243    async fn sleeping_run(store: &InMemoryStore, scheduled_at: DateTime<Utc>) -> Run {
4244        let run = store
4245            .create_run(new_run_req("sleepy"))
4246            .await
4247            .unwrap()
4248            .into_run();
4249        store
4250            .update_run_status(run.id, RunStatus::Running)
4251            .await
4252            .unwrap();
4253        store
4254            .update_run(
4255                run.id,
4256                RunUpdate {
4257                    status: Some(RunStatus::Sleeping),
4258                    scheduled_at: Some(scheduled_at),
4259                    ..RunUpdate::default()
4260                },
4261            )
4262            .await
4263            .unwrap();
4264        store.get_run(run.id).await.unwrap().unwrap()
4265    }
4266
4267    #[tokio::test]
4268    async fn claim_due_sleeping_runs_requeues_due_runs() {
4269        let store = InMemoryStore::new();
4270        let due = sleeping_run(&store, Utc::now() - TimeDelta::seconds(5)).await;
4271
4272        let woken = store.claim_due_sleeping_runs(10).await.unwrap();
4273        assert_eq!(woken.len(), 1);
4274        assert_eq!(woken[0].id, due.id);
4275        assert_eq!(woken[0].status.state, RunStatus::Pending);
4276        assert!(woken[0].scheduled_at.is_none());
4277
4278        let fetched = store.get_run(due.id).await.unwrap().unwrap();
4279        assert_eq!(fetched.status.state, RunStatus::Pending);
4280        assert!(fetched.scheduled_at.is_none());
4281
4282        // Exactly once: a second tick finds nothing.
4283        assert!(store.claim_due_sleeping_runs(10).await.unwrap().is_empty());
4284    }
4285
4286    #[tokio::test]
4287    async fn claim_due_sleeping_runs_skips_future_runs() {
4288        let store = InMemoryStore::new();
4289        let future = sleeping_run(&store, Utc::now() + TimeDelta::hours(1)).await;
4290
4291        assert!(store.claim_due_sleeping_runs(10).await.unwrap().is_empty());
4292        let fetched = store.get_run(future.id).await.unwrap().unwrap();
4293        assert_eq!(fetched.status.state, RunStatus::Sleeping);
4294    }
4295
4296    #[tokio::test]
4297    async fn claim_due_sleeping_runs_honours_limit_oldest_first() {
4298        let store = InMemoryStore::new();
4299        let older = sleeping_run(&store, Utc::now() - TimeDelta::seconds(20)).await;
4300        let newer = sleeping_run(&store, Utc::now() - TimeDelta::seconds(10)).await;
4301
4302        let woken = store.claim_due_sleeping_runs(1).await.unwrap();
4303        assert_eq!(woken.len(), 1);
4304        assert_eq!(woken[0].id, older.id);
4305
4306        let woken = store.claim_due_sleeping_runs(1).await.unwrap();
4307        assert_eq!(woken[0].id, newer.id);
4308    }
4309}