Skip to main content

ironflow_store/memory/
run_store.rs

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