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