Skip to main content

ironflow_store/memory/
run_store.rs

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