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