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