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