Skip to main content

ironflow_store/memory/
run_store.rs

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