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