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