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