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