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