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