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