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