1use std::collections::HashMap;
10use std::fmt;
11use std::sync::Arc;
12use std::time::Instant;
13
14use chrono::{DateTime, TimeDelta, Utc};
15use rust_decimal::Decimal;
16use serde_json::Value;
17use tracing::{error, info, warn};
18use uuid::Uuid;
19
20#[cfg(feature = "prometheus")]
21use ironflow_core::metric_names::{
22 RUN_BUDGET_EXCEEDED_TOTAL, RUN_COST_USD, RUN_DURATION_SECONDS, RUNS_ACTIVE, RUNS_TOTAL,
23};
24use ironflow_core::provider::AgentProvider;
25use ironflow_store::error::StoreError;
26use ironflow_store::models::{
27 NewRun, Run, RunActor, RunCreation, RunFilter, RunStatus, RunUpdate, StepStatus, StepUpdate,
28 TriggerKind,
29};
30use ironflow_store::store::Store;
31#[cfg(feature = "prometheus")]
32use metrics::{counter, gauge, histogram};
33
34use crate::budget::{BudgetConfig, month_start};
35use crate::context::WorkflowContext;
36use crate::error::EngineError;
37use crate::handler::{WorkflowHandler, WorkflowInfo};
38use crate::log_sender::LogSender;
39use crate::notify::{Event, EventPublisher, EventSubscriber};
40use crate::retry_policy::{backoff_for_retry, is_run_retryable};
41use crate::schedule::CronSchedule;
42
43#[derive(Debug, Clone, Default)]
62pub struct EnqueueOptions {
63 pub max_retries: u32,
65 pub labels: HashMap<String, String>,
67 pub scheduled_at: Option<DateTime<Utc>>,
70 pub max_cost_usd: Option<Decimal>,
74 pub created_by: Option<RunActor>,
77 pub idempotency_key: Option<String>,
83}
84
85pub struct Engine {
126 store: Arc<dyn Store>,
127 provider: Arc<dyn AgentProvider>,
128 handlers: HashMap<String, Arc<dyn WorkflowHandler>>,
129 event_publisher: EventPublisher,
130 log_sender: Option<LogSender>,
131 budget: BudgetConfig,
132}
133
134fn validate_category(handler_name: &str, category: &str) -> Result<(), EngineError> {
144 let reject = |reason: &str| {
145 Err(EngineError::InvalidWorkflow(format!(
146 "handler '{handler_name}' has invalid category '{category}': {reason}"
147 )))
148 };
149
150 if category.is_empty() {
151 return reject("empty category");
152 }
153 if category.starts_with('/') {
154 return reject("leading '/'");
155 }
156 if category.ends_with('/') {
157 return reject("trailing '/'");
158 }
159 for segment in category.split('/') {
160 if segment.is_empty() {
161 return reject("empty segment (double '/')");
162 }
163 if segment.trim().is_empty() {
164 return reject("whitespace-only segment");
165 }
166 }
167 Ok(())
168}
169
170impl Engine {
171 pub fn new(store: Arc<dyn Store>, provider: Arc<dyn AgentProvider>) -> Self {
187 Self {
188 store,
189 provider,
190 handlers: HashMap::new(),
191 event_publisher: EventPublisher::new(),
192 log_sender: None,
193 budget: BudgetConfig::new(),
194 }
195 }
196
197 pub fn with_budget_config(mut self, budget: BudgetConfig) -> Self {
218 self.budget = budget;
219 self
220 }
221
222 pub fn budget_config(&self) -> &BudgetConfig {
224 &self.budget
225 }
226
227 pub fn set_log_sender(&mut self, sender: LogSender) {
233 self.log_sender = Some(sender);
234 }
235
236 pub fn store(&self) -> &Arc<dyn Store> {
238 &self.store
239 }
240
241 pub fn provider(&self) -> &Arc<dyn AgentProvider> {
243 &self.provider
244 }
245
246 fn build_context(&self, run: &Run) -> WorkflowContext {
255 let handlers = self.handlers.clone();
256 let resolver: crate::context::HandlerResolver =
257 Arc::new(move |name: &str| handlers.get(name).cloned());
258 let mut ctx = WorkflowContext::with_handler_resolver(
259 run.id,
260 self.store.clone(),
261 self.provider.clone(),
262 resolver,
263 );
264 ctx.carry_over_run_totals(run.retry_count + 1, run.cost_usd, run.duration_ms);
265 ctx.set_max_cost_usd(run.max_cost_usd);
266 if let Some(ref sender) = self.log_sender {
267 ctx.set_log_sender(sender.clone());
268 }
269 ctx
270 }
271
272 async fn check_monthly_quota(&self, workflow_name: &str) -> Result<(), EngineError> {
283 let Some(limit) = self.budget.monthly_cost_limit_usd else {
284 return Ok(());
285 };
286
287 let stats = self
288 .store
289 .get_stats(RunFilter {
290 created_after: Some(month_start(Utc::now())),
291 ..RunFilter::default()
292 })
293 .await?;
294
295 if stats.total_cost_usd < limit {
296 return Ok(());
297 }
298
299 warn!(
300 workflow = %workflow_name,
301 limit_usd = %limit,
302 spent_usd = %stats.total_cost_usd,
303 "monthly cost quota exhausted, refusing new run"
304 );
305
306 #[cfg(feature = "prometheus")]
307 counter!(
308 RUN_BUDGET_EXCEEDED_TOTAL,
309 "workflow" => workflow_name.to_string(),
310 "scope" => "monthly",
311 )
312 .increment(1);
313
314 Err(EngineError::MonthlyBudgetExceeded {
315 limit_usd: limit,
316 spent_usd: stats.total_cost_usd,
317 })
318 }
319
320 pub fn register(&mut self, handler: impl WorkflowHandler + 'static) -> Result<(), EngineError> {
364 let name = handler.name().to_string();
365 if self.handlers.contains_key(&name) {
366 return Err(EngineError::InvalidWorkflow(format!(
367 "handler '{}' already registered",
368 name
369 )));
370 }
371 if let Some(category) = handler.category() {
372 validate_category(&name, category)?;
373 }
374 self.handlers.insert(name, Arc::new(handler));
375 Ok(())
376 }
377
378 pub fn register_boxed(&mut self, handler: Box<dyn WorkflowHandler>) -> Result<(), EngineError> {
385 let name = handler.name().to_string();
386 if self.handlers.contains_key(&name) {
387 return Err(EngineError::InvalidWorkflow(format!(
388 "handler '{}' already registered",
389 name
390 )));
391 }
392 if let Some(category) = handler.category() {
393 validate_category(&name, category)?;
394 }
395 self.handlers.insert(name, Arc::from(handler));
396 Ok(())
397 }
398
399 pub fn get_handler(&self, name: &str) -> Option<&Arc<dyn WorkflowHandler>> {
401 self.handlers.get(name)
402 }
403
404 pub fn handler_names(&self) -> Vec<&str> {
406 self.handlers.keys().map(|s| s.as_str()).collect()
407 }
408
409 pub fn handler_info(&self, name: &str) -> Option<WorkflowInfo> {
411 self.handlers.get(name).map(|h| h.describe())
412 }
413
414 pub fn scheduled_handlers(&self) -> Vec<(&str, &CronSchedule)> {
438 self.handlers
439 .iter()
440 .filter_map(|(name, handler)| handler.schedule().map(|sched| (name.as_str(), sched)))
441 .collect()
442 }
443
444 pub fn subscribe(
469 &mut self,
470 subscriber: impl EventSubscriber + 'static,
471 event_types: &[&'static str],
472 ) {
473 self.event_publisher.subscribe(subscriber, event_types);
474 }
475
476 pub fn event_publisher(&self) -> &EventPublisher {
481 &self.event_publisher
482 }
483
484 #[tracing::instrument(name = "engine.run_handler", skip_all, fields(workflow = %handler_name))]
514 pub async fn run_handler(
515 &self,
516 handler_name: &str,
517 trigger: TriggerKind,
518 payload: Value,
519 ) -> Result<Run, EngineError> {
520 let handler = self
521 .handlers
522 .get(handler_name)
523 .ok_or_else(|| {
524 EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
525 })?
526 .clone();
527
528 self.check_monthly_quota(handler_name).await?;
529
530 let handler_version = handler.version().map(str::to_string);
531 let max_cost_usd = self
532 .budget
533 .resolve_run_cap(None, handler.default_max_cost_usd());
534 let run = self
535 .store
536 .create_run(NewRun {
537 created_by: None,
538 workflow_name: handler_name.to_string(),
539 trigger,
540 payload,
541 max_retries: 0,
542 handler_version,
543 labels: handler.default_labels(),
544 scheduled_at: None,
545 idempotency_key: None,
546 max_cost_usd,
547 })
548 .await?
549 .into_run();
550
551 let run_id = run.id;
552 info!(run_id = %run_id, handler_version = run.handler_version.as_deref().unwrap_or(""), "run created");
553
554 self.store
555 .update_run_status(run_id, RunStatus::Running)
556 .await?;
557
558 #[cfg(feature = "prometheus")]
559 gauge!(RUNS_ACTIVE, "workflow" => handler_name.to_string()).increment(1.0);
560
561 let run_start = Instant::now();
562 let mut ctx = self.build_context(&run);
563
564 let result = handler.execute(&mut ctx).await;
565 self.finalize_run(run_id, handler_name, result, &ctx, run_start)
566 .await
567 }
568
569 #[tracing::instrument(name = "engine.enqueue_handler", skip_all, fields(workflow = %handler_name))]
580 pub async fn enqueue_handler(
581 &self,
582 handler_name: &str,
583 trigger: TriggerKind,
584 payload: Value,
585 max_retries: u32,
586 ) -> Result<Run, EngineError> {
587 self.enqueue_handler_with_options(
588 handler_name,
589 trigger,
590 payload,
591 EnqueueOptions {
592 max_retries,
593 ..Default::default()
594 },
595 )
596 .await
597 .map(RunCreation::into_run)
598 }
599
600 #[tracing::instrument(name = "engine.enqueue_handler_with_options", skip_all, fields(workflow = %handler_name))]
645 pub async fn enqueue_handler_with_options(
646 &self,
647 handler_name: &str,
648 trigger: TriggerKind,
649 payload: Value,
650 options: EnqueueOptions,
651 ) -> Result<RunCreation, EngineError> {
652 let EnqueueOptions {
653 max_retries,
654 labels,
655 scheduled_at,
656 max_cost_usd,
657 created_by,
658 idempotency_key,
659 } = options;
660
661 let handler = self.handlers.get(handler_name).ok_or_else(|| {
662 EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
663 })?;
664
665 self.check_monthly_quota(handler_name).await?;
666
667 let handler_version = handler.version().map(str::to_string);
668 let mut merged_labels = handler.default_labels();
669 merged_labels.extend(labels);
670 let resolved_cap = self
671 .budget
672 .resolve_run_cap(max_cost_usd, handler.default_max_cost_usd());
673
674 let creation = self
675 .store
676 .create_run(NewRun {
677 workflow_name: handler_name.to_string(),
678 trigger,
679 payload,
680 max_retries,
681 handler_version,
682 labels: merged_labels,
683 scheduled_at,
684 created_by,
685 idempotency_key,
686 max_cost_usd: resolved_cap,
687 })
688 .await?;
689
690 match &creation {
691 RunCreation::Created(run) => info!(
692 run_id = %run.id,
693 workflow = %handler_name,
694 max_cost_usd = ?resolved_cap,
695 "handler run enqueued"
696 ),
697 RunCreation::Existing(run) => info!(
698 run_id = %run.id,
699 workflow = %handler_name,
700 "idempotent replay, nothing enqueued"
701 ),
702 }
703
704 Ok(creation)
705 }
706
707 #[tracing::instrument(name = "engine.execute_handler_run", skip_all, fields(run_id = %run_id))]
716 pub async fn execute_handler_run(&self, run_id: Uuid) -> Result<Run, EngineError> {
717 let run = self
718 .store
719 .get_run(run_id)
720 .await?
721 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
722
723 let handler = self
724 .handlers
725 .get(&run.workflow_name)
726 .ok_or_else(|| {
727 EngineError::InvalidWorkflow(format!(
728 "no handler registered: {}",
729 run.workflow_name
730 ))
731 })?
732 .clone();
733
734 #[cfg(feature = "prometheus")]
735 gauge!(RUNS_ACTIVE, "workflow" => run.workflow_name.clone()).increment(1.0);
736
737 let run_start = Instant::now();
738 let mut ctx = self.build_context(&run);
739
740 if run.retry_count > 0 {
743 ctx.load_replay_steps().await?;
744 }
745
746 let result = handler.execute(&mut ctx).await;
747 self.finalize_run(run_id, &run.workflow_name, result, &ctx, run_start)
748 .await
749 }
750
751 #[tracing::instrument(name = "engine.execute_run", skip_all, fields(run_id = %run_id))]
759 pub async fn execute_run(&self, run_id: Uuid) -> Result<Run, EngineError> {
760 self.execute_handler_run(run_id).await
761 }
762
763 #[tracing::instrument(name = "engine.resume_run", skip_all, fields(run_id = %run_id))]
777 pub async fn resume_run(&self, run_id: Uuid) -> Result<Run, EngineError> {
778 let run = self
779 .store
780 .get_run(run_id)
781 .await?
782 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
783
784 let handler = self
785 .handlers
786 .get(&run.workflow_name)
787 .ok_or_else(|| {
788 EngineError::InvalidWorkflow(format!(
789 "no handler registered: {}",
790 run.workflow_name
791 ))
792 })?
793 .clone();
794
795 info!(run_id = %run_id, workflow = %run.workflow_name, "resuming run after approval");
796
797 let run_start = Instant::now();
798 let mut ctx = self.build_context(&run);
799 ctx.load_replay_steps().await?;
800
801 let result = handler.execute(&mut ctx).await;
802 self.finalize_run(run_id, &run.workflow_name, result, &ctx, run_start)
803 .await
804 }
805
806 pub async fn fail_or_schedule_retry(
851 &self,
852 run_id: Uuid,
853 error: &str,
854 retryable: bool,
855 cost_usd: Option<Decimal>,
856 duration_ms: Option<u64>,
857 ) -> Result<RunStatus, EngineError> {
858 let run = self
859 .store
860 .get_run(run_id)
861 .await?
862 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
863
864 let has_attempts_left = run.retry_count < run.max_retries;
865 let update = if retryable && has_attempts_left {
866 let backoff = backoff_for_retry(run.retry_count);
867 let scheduled_at = Utc::now() + TimeDelta::milliseconds(backoff.as_millis() as i64);
868
869 info!(
870 run_id = %run_id,
871 workflow = %run.workflow_name,
872 attempt = run.retry_count + 1,
873 max_retries = run.max_retries,
874 backoff_secs = backoff.as_secs(),
875 scheduled_at = %scheduled_at,
876 "run failed, scheduling retry"
877 );
878
879 RunUpdate {
880 status: Some(RunStatus::Retrying),
881 error: Some(error.to_string()),
882 increment_retry: true,
883 cost_usd,
884 duration_ms,
885 scheduled_at: Some(scheduled_at),
886 ..RunUpdate::default()
887 }
888 } else {
889 RunUpdate {
890 status: Some(RunStatus::Failed),
891 error: Some(error.to_string()),
892 cost_usd,
893 duration_ms,
894 completed_at: Some(Utc::now()),
895 ..RunUpdate::default()
896 }
897 };
898
899 let status = update.status.unwrap_or(RunStatus::Failed);
900 self.store.update_run(run_id, update).await?;
901 self.fail_orphaned_steps(run_id, error).await?;
902
903 Ok(status)
904 }
905
906 pub async fn fail_orphaned_steps(
920 &self,
921 run_id: Uuid,
922 error_message: &str,
923 ) -> Result<(), EngineError> {
924 let steps = self.store.list_steps(run_id).await?;
925 let now = Utc::now();
926
927 for step in steps {
928 if step.status.state.is_terminal() {
929 continue;
930 }
931
932 let (target_status, error) = match step.status.state {
933 StepStatus::Running | StepStatus::AwaitingApproval => {
934 let err = if step.error.is_some() {
935 None
936 } else {
937 Some(error_message.to_string())
938 };
939 (StepStatus::Failed, err)
940 }
941 StepStatus::Pending => (StepStatus::Skipped, None),
942 _ => continue,
943 };
944
945 if let Err(e) = self
946 .store
947 .update_step(
948 step.id,
949 StepUpdate {
950 status: Some(target_status),
951 error,
952 completed_at: Some(now),
953 ..StepUpdate::default()
954 },
955 )
956 .await
957 {
958 warn!(
959 run_id = %run_id,
960 step_id = %step.id,
961 step_name = %step.name,
962 error = %e,
963 "failed to cleanup orphaned step"
964 );
965 } else {
966 info!(
967 run_id = %run_id,
968 step_id = %step.id,
969 step_name = %step.name,
970 from = %step.status.state,
971 to = %target_status,
972 "cleaned up orphaned step"
973 );
974 }
975 }
976
977 Ok(())
978 }
979
980 async fn finalize_run(
986 &self,
987 run_id: Uuid,
988 workflow_name: &str,
989 result: Result<(), EngineError>,
990 ctx: &WorkflowContext,
991 run_start: Instant,
992 ) -> Result<Run, EngineError> {
993 let total_duration = ctx.carried_duration_ms() + run_start.elapsed().as_millis() as u64;
996 let completed_at = Utc::now();
997
998 let final_status;
999 let final_run;
1000
1001 match result {
1002 Ok(()) => {
1003 final_status = RunStatus::Completed;
1004 final_run = self
1005 .store
1006 .update_run_returning(
1007 run_id,
1008 RunUpdate {
1009 status: Some(RunStatus::Completed),
1010 cost_usd: Some(ctx.total_cost_usd()),
1011 duration_ms: Some(total_duration),
1012 completed_at: Some(completed_at),
1013 ..RunUpdate::default()
1014 },
1015 )
1016 .await?;
1017
1018 info!(
1019 run_id = %run_id,
1020 cost_usd = %ctx.total_cost_usd(),
1021 duration_ms = total_duration,
1022 "run completed"
1023 );
1024 }
1025 Err(EngineError::ApprovalRequired {
1026 run_id: approval_run_id,
1027 step_id,
1028 ref message,
1029 }) => {
1030 final_status = RunStatus::AwaitingApproval;
1031 final_run = self
1032 .store
1033 .update_run_returning(
1034 run_id,
1035 RunUpdate {
1036 status: Some(RunStatus::AwaitingApproval),
1037 cost_usd: Some(ctx.total_cost_usd()),
1038 duration_ms: Some(total_duration),
1039 ..RunUpdate::default()
1040 },
1041 )
1042 .await?;
1043
1044 info!(
1045 run_id = %approval_run_id,
1046 step_id = %step_id,
1047 message = %message,
1048 "run awaiting approval"
1049 );
1050 }
1051 Err(err) => {
1052 let budget_exceeded = matches!(err, EngineError::RunBudgetExceeded { .. });
1056
1057 final_status = if budget_exceeded {
1058 if let Err(store_err) = self
1059 .store
1060 .update_run(
1061 run_id,
1062 RunUpdate {
1063 status: Some(RunStatus::Cancelled),
1064 error: Some(err.to_string()),
1065 cost_usd: Some(ctx.total_cost_usd()),
1066 duration_ms: Some(total_duration),
1067 completed_at: Some(completed_at),
1068 ..RunUpdate::default()
1069 },
1070 )
1071 .await
1072 {
1073 error!(run_id = %run_id, store_error = %store_err, "failed to persist run cancellation");
1074 }
1075 if let Err(cleanup_err) = self
1076 .fail_orphaned_steps(run_id, "run stopped: cost cap reached")
1077 .await
1078 {
1079 error!(run_id = %run_id, store_error = %cleanup_err, "failed to cleanup orphaned steps");
1080 }
1081 RunStatus::Cancelled
1082 } else {
1083 self.fail_or_schedule_retry(
1084 run_id,
1085 &err.to_string(),
1086 is_run_retryable(&err),
1087 Some(ctx.total_cost_usd()),
1088 Some(total_duration),
1089 )
1090 .await
1091 .unwrap_or_else(|store_err| {
1092 error!(run_id = %run_id, store_error = %store_err, "failed to persist run failure");
1093 RunStatus::Failed
1094 })
1095 };
1096
1097 if budget_exceeded {
1098 self.on_run_budget_exceeded(workflow_name, run_id, &err);
1099 }
1100
1101 error!(run_id = %run_id, status = %final_status, error = %err, "run stopped");
1102
1103 self.publish_run_status_changed(
1104 workflow_name,
1105 run_id,
1106 final_status,
1107 Some(err.to_string()),
1108 ctx,
1109 total_duration,
1110 );
1111
1112 #[cfg(feature = "prometheus")]
1113 self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
1114
1115 return Err(err);
1116 }
1117 }
1118
1119 self.publish_run_status_changed(
1120 workflow_name,
1121 run_id,
1122 final_status,
1123 None,
1124 ctx,
1125 total_duration,
1126 );
1127
1128 #[cfg(feature = "prometheus")]
1129 self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
1130
1131 Ok(final_run)
1132 }
1133
1134 #[cfg(feature = "prometheus")]
1136 fn emit_run_metrics(
1137 &self,
1138 workflow_name: &str,
1139 status: RunStatus,
1140 duration_ms: u64,
1141 ctx: &WorkflowContext,
1142 ) {
1143 let status_str = status.to_string();
1144 let wf = workflow_name.to_string();
1145
1146 counter!(RUNS_TOTAL, "workflow" => wf.clone(), "status" => status_str.clone()).increment(1);
1147 histogram!(RUN_DURATION_SECONDS, "workflow" => wf.clone(), "status" => status_str)
1148 .record(duration_ms as f64 / 1000.0);
1149 histogram!(RUN_COST_USD, "workflow" => wf.clone()).record(
1150 ctx.total_cost_usd()
1151 .to_string()
1152 .parse::<f64>()
1153 .unwrap_or(0.0),
1154 );
1155 gauge!(RUNS_ACTIVE, "workflow" => wf).decrement(1.0);
1156 }
1157
1158 fn on_run_budget_exceeded(&self, workflow_name: &str, run_id: Uuid, err: &EngineError) {
1164 let EngineError::RunBudgetExceeded {
1165 limit_usd,
1166 spent_usd,
1167 step_budget_usd,
1168 ..
1169 } = err
1170 else {
1171 return;
1172 };
1173
1174 #[cfg(feature = "prometheus")]
1175 counter!(
1176 RUN_BUDGET_EXCEEDED_TOTAL,
1177 "workflow" => workflow_name.to_string(),
1178 "scope" => "run",
1179 )
1180 .increment(1);
1181
1182 self.event_publisher.publish(Event::RunBudgetExceeded {
1183 run_id,
1184 workflow_name: workflow_name.to_string(),
1185 limit_usd: *limit_usd,
1186 spent_usd: *spent_usd,
1187 step_budget_usd: *step_budget_usd,
1188 at: Utc::now(),
1189 });
1190 }
1191
1192 fn publish_run_status_changed(
1197 &self,
1198 workflow_name: &str,
1199 run_id: Uuid,
1200 to: RunStatus,
1201 error: Option<String>,
1202 ctx: &WorkflowContext,
1203 duration_ms: u64,
1204 ) {
1205 let now = Utc::now();
1206 let cost_usd = ctx.total_cost_usd();
1207 let wf = workflow_name.to_string();
1208
1209 self.event_publisher.publish(Event::RunStatusChanged {
1210 run_id,
1211 workflow_name: wf.clone(),
1212 from: RunStatus::Running,
1213 to,
1214 error: error.clone(),
1215 cost_usd,
1216 duration_ms,
1217 at: now,
1218 });
1219
1220 if to == RunStatus::Failed {
1221 self.event_publisher.publish(Event::RunFailed {
1222 run_id,
1223 workflow_name: wf,
1224 error,
1225 cost_usd,
1226 duration_ms,
1227 at: now,
1228 });
1229 }
1230 }
1231}
1232
1233impl fmt::Debug for Engine {
1234 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1235 f.debug_struct("Engine")
1236 .field("handlers", &self.handlers.keys().collect::<Vec<_>>())
1237 .finish_non_exhaustive()
1238 }
1239}
1240
1241#[cfg(test)]
1242mod tests {
1243 use super::*;
1244 use crate::config::ShellConfig;
1245 use crate::handler::{HandlerFuture, WorkflowHandler};
1246 use ironflow_core::providers::claude::ClaudeCodeProvider;
1247 use ironflow_core::providers::record_replay::RecordReplayProvider;
1248 use ironflow_store::memory::InMemoryStore;
1249 use ironflow_store::models::StepStatus;
1250 use serde_json::json;
1251
1252 struct EchoWorkflow;
1254
1255 impl WorkflowHandler for EchoWorkflow {
1256 fn name(&self) -> &str {
1257 "echo-workflow"
1258 }
1259
1260 fn describe(&self) -> WorkflowInfo {
1261 WorkflowInfo {
1262 description: "A simple workflow that echoes hello".to_string(),
1263 source_code: None,
1264 sub_workflows: Vec::new(),
1265 category: None,
1266 version: self.version().map(str::to_string),
1267 input_schema: None,
1268 default_labels: HashMap::new(),
1269 schedule: self.schedule().cloned(),
1270 default_max_cost_usd: self.default_max_cost_usd(),
1271 }
1272 }
1273
1274 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
1275 Box::pin(async move {
1276 ctx.shell("greet", ShellConfig::new("echo hello")).await?;
1277 Ok(())
1278 })
1279 }
1280 }
1281
1282 struct FailingWorkflow;
1284
1285 impl WorkflowHandler for FailingWorkflow {
1286 fn name(&self) -> &str {
1287 "failing-workflow"
1288 }
1289
1290 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
1291 Box::pin(async move {
1292 ctx.shell("fail", ShellConfig::new("exit 1")).await?;
1293 Ok(())
1294 })
1295 }
1296 }
1297
1298 fn create_test_engine() -> Engine {
1299 let store = Arc::new(InMemoryStore::new());
1300 let inner = ClaudeCodeProvider::new();
1301 let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
1302 inner,
1303 "/tmp/ironflow-fixtures",
1304 ));
1305 Engine::new(store, provider)
1306 }
1307
1308 #[test]
1309 fn engine_new_creates_instance() {
1310 let engine = create_test_engine();
1311 assert_eq!(engine.handler_names().len(), 0);
1312 }
1313
1314 #[test]
1315 fn engine_register_handler() {
1316 let mut engine = create_test_engine();
1317 let result = engine.register(EchoWorkflow);
1318 assert!(result.is_ok());
1319 assert_eq!(engine.handler_names().len(), 1);
1320 assert!(engine.handler_names().contains(&"echo-workflow"));
1321 }
1322
1323 #[test]
1324 fn engine_register_duplicate_returns_error() {
1325 let mut engine = create_test_engine();
1326 engine.register(EchoWorkflow).unwrap();
1327 let result = engine.register(EchoWorkflow);
1328 assert!(result.is_err());
1329 }
1330
1331 #[test]
1332 fn engine_get_handler_found() {
1333 let mut engine = create_test_engine();
1334 engine.register(EchoWorkflow).unwrap();
1335 let handler = engine.get_handler("echo-workflow");
1336 assert!(handler.is_some());
1337 }
1338
1339 #[test]
1340 fn engine_get_handler_not_found() {
1341 let engine = create_test_engine();
1342 let handler = engine.get_handler("nonexistent");
1343 assert!(handler.is_none());
1344 }
1345
1346 #[test]
1347 fn engine_handler_names_lists_all() {
1348 let mut engine = create_test_engine();
1349 engine.register(EchoWorkflow).unwrap();
1350 engine.register(FailingWorkflow).unwrap();
1351 let names = engine.handler_names();
1352 assert_eq!(names.len(), 2);
1353 assert!(names.contains(&"echo-workflow"));
1354 assert!(names.contains(&"failing-workflow"));
1355 }
1356
1357 #[test]
1358 fn engine_handler_info_returns_description() {
1359 let mut engine = create_test_engine();
1360 engine.register(EchoWorkflow).unwrap();
1361 let info = engine.handler_info("echo-workflow");
1362 assert!(info.is_some());
1363 let info = info.unwrap();
1364 assert_eq!(info.description, "A simple workflow that echoes hello");
1365 }
1366
1367 struct CategorizedWorkflow;
1368
1369 impl WorkflowHandler for CategorizedWorkflow {
1370 fn name(&self) -> &str {
1371 "categorized"
1372 }
1373 fn category(&self) -> Option<&str> {
1374 Some("data/etl")
1375 }
1376 fn execute<'a>(
1377 &'a self,
1378 _ctx: &'a mut WorkflowContext,
1379 ) -> crate::handler::HandlerFuture<'a> {
1380 Box::pin(async move { Ok(()) })
1381 }
1382 }
1383
1384 #[test]
1385 fn engine_default_describe_propagates_category() {
1386 let mut engine = create_test_engine();
1387 engine.register(CategorizedWorkflow).unwrap();
1388 let info = engine.handler_info("categorized").unwrap();
1389 assert_eq!(info.category.as_deref(), Some("data/etl"));
1390 }
1391
1392 #[test]
1393 fn engine_default_describe_without_category() {
1394 let mut engine = create_test_engine();
1395 engine.register(EchoWorkflow).unwrap();
1396 let info = engine.handler_info("echo-workflow").unwrap();
1397 assert!(info.category.is_none());
1398 }
1399
1400 struct ScheduledWorkflow {
1405 schedule: CronSchedule,
1406 }
1407
1408 impl ScheduledWorkflow {
1409 fn new() -> Self {
1410 Self {
1411 schedule: CronSchedule::new("0 0 * * * *").unwrap(),
1412 }
1413 }
1414 }
1415
1416 impl WorkflowHandler for ScheduledWorkflow {
1417 fn name(&self) -> &str {
1418 "scheduled"
1419 }
1420 fn schedule(&self) -> Option<&CronSchedule> {
1421 Some(&self.schedule)
1422 }
1423 fn execute<'a>(
1424 &'a self,
1425 _ctx: &'a mut WorkflowContext,
1426 ) -> crate::handler::HandlerFuture<'a> {
1427 Box::pin(async move { Ok(()) })
1428 }
1429 }
1430
1431 #[test]
1432 fn engine_default_describe_propagates_schedule() {
1433 let mut engine = create_test_engine();
1434 engine.register(ScheduledWorkflow::new()).unwrap();
1435 let info = engine.handler_info("scheduled").unwrap();
1436 assert_eq!(
1437 info.schedule.as_ref().map(|s| s.as_str()),
1438 Some("0 0 * * * *")
1439 );
1440 }
1441
1442 #[test]
1443 fn engine_default_describe_without_schedule() {
1444 let mut engine = create_test_engine();
1445 engine.register(EchoWorkflow).unwrap();
1446 let info = engine.handler_info("echo-workflow").unwrap();
1447 assert!(info.schedule.is_none());
1448 }
1449
1450 #[test]
1451 fn scheduled_handlers_returns_only_scheduled() {
1452 let mut engine = create_test_engine();
1453 engine.register(EchoWorkflow).unwrap();
1454 engine.register(ScheduledWorkflow::new()).unwrap();
1455 engine.register(FailingWorkflow).unwrap();
1456
1457 let scheduled = engine.scheduled_handlers();
1458 assert_eq!(scheduled.len(), 1);
1459 assert_eq!(scheduled[0].0, "scheduled");
1460 assert_eq!(scheduled[0].1.as_str(), "0 0 * * * *");
1461 }
1462
1463 #[test]
1464 fn scheduled_handlers_empty_when_none_scheduled() {
1465 let mut engine = create_test_engine();
1466 engine.register(EchoWorkflow).unwrap();
1467 engine.register(FailingWorkflow).unwrap();
1468
1469 let scheduled = engine.scheduled_handlers();
1470 assert!(scheduled.is_empty());
1471 }
1472
1473 struct BadCategoryWorkflow(&'static str);
1474
1475 impl WorkflowHandler for BadCategoryWorkflow {
1476 fn name(&self) -> &str {
1477 "bad-category"
1478 }
1479 fn category(&self) -> Option<&str> {
1480 Some(self.0)
1481 }
1482 fn execute<'a>(
1483 &'a self,
1484 _ctx: &'a mut WorkflowContext,
1485 ) -> crate::handler::HandlerFuture<'a> {
1486 Box::pin(async move { Ok(()) })
1487 }
1488 }
1489
1490 #[test]
1491 fn engine_register_rejects_empty_category() {
1492 let mut engine = create_test_engine();
1493 let err = engine.register(BadCategoryWorkflow("")).unwrap_err();
1494 match err {
1495 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty category")),
1496 other => panic!("expected InvalidWorkflow, got {other:?}"),
1497 }
1498 }
1499
1500 #[test]
1501 fn engine_register_rejects_leading_slash_category() {
1502 let mut engine = create_test_engine();
1503 let err = engine
1504 .register(BadCategoryWorkflow("/data/etl"))
1505 .unwrap_err();
1506 match err {
1507 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("leading '/'")),
1508 other => panic!("expected InvalidWorkflow, got {other:?}"),
1509 }
1510 }
1511
1512 #[test]
1513 fn engine_register_rejects_trailing_slash_category() {
1514 let mut engine = create_test_engine();
1515 let err = engine
1516 .register(BadCategoryWorkflow("data/etl/"))
1517 .unwrap_err();
1518 match err {
1519 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("trailing '/'")),
1520 other => panic!("expected InvalidWorkflow, got {other:?}"),
1521 }
1522 }
1523
1524 #[test]
1525 fn engine_register_rejects_double_slash_category() {
1526 let mut engine = create_test_engine();
1527 let err = engine
1528 .register(BadCategoryWorkflow("data//etl"))
1529 .unwrap_err();
1530 match err {
1531 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty segment")),
1532 other => panic!("expected InvalidWorkflow, got {other:?}"),
1533 }
1534 }
1535
1536 #[test]
1537 fn engine_register_rejects_whitespace_only_segment_category() {
1538 let mut engine = create_test_engine();
1539 let err = engine
1540 .register(BadCategoryWorkflow("data/ /etl"))
1541 .unwrap_err();
1542 match err {
1543 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("whitespace-only segment")),
1544 other => panic!("expected InvalidWorkflow, got {other:?}"),
1545 }
1546 }
1547
1548 #[test]
1549 fn engine_register_accepts_valid_nested_category() {
1550 let mut engine = create_test_engine();
1551 assert!(engine.register(CategorizedWorkflow).is_ok());
1552 }
1553
1554 #[tokio::test]
1555 async fn engine_unknown_workflow_returns_error() {
1556 let engine = create_test_engine();
1557 let result = engine
1558 .run_handler("unknown", TriggerKind::Manual, json!({}))
1559 .await;
1560 assert!(result.is_err());
1561 match result {
1562 Err(EngineError::InvalidWorkflow(msg)) => {
1563 assert!(msg.contains("no handler registered"));
1564 }
1565 _ => panic!("expected InvalidWorkflow error"),
1566 }
1567 }
1568
1569 #[tokio::test]
1570 async fn engine_enqueue_handler_creates_pending_run() {
1571 let mut engine = create_test_engine();
1572 engine.register(EchoWorkflow).unwrap();
1573
1574 let run = engine
1575 .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
1576 .await
1577 .unwrap();
1578 assert_eq!(run.status.state, RunStatus::Pending);
1579 assert_eq!(run.workflow_name, "echo-workflow");
1580 }
1581
1582 #[tokio::test]
1583 async fn enqueue_handler_leaves_the_run_unattributed() {
1584 let mut engine = create_test_engine();
1585 engine.register(EchoWorkflow).unwrap();
1586
1587 let run = engine
1588 .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
1589 .await
1590 .unwrap();
1591
1592 assert!(run.created_by.is_none());
1593 }
1594
1595 #[tokio::test]
1596 async fn enqueue_handler_with_options_records_the_author() {
1597 let mut engine = create_test_engine();
1598 engine.register(EchoWorkflow).unwrap();
1599 let actor = RunActor::User {
1600 user_id: Uuid::now_v7(),
1601 };
1602
1603 let run = engine
1604 .enqueue_handler_with_options(
1605 "echo-workflow",
1606 TriggerKind::Api,
1607 json!({}),
1608 EnqueueOptions {
1609 created_by: Some(actor.clone()),
1610 ..Default::default()
1611 },
1612 )
1613 .await
1614 .unwrap()
1615 .into_run();
1616
1617 assert_eq!(run.created_by, Some(actor));
1618 }
1619
1620 #[tokio::test]
1621 async fn enqueue_handler_with_options_accepts_no_author() {
1622 let mut engine = create_test_engine();
1623 engine.register(EchoWorkflow).unwrap();
1624
1625 let run = engine
1626 .enqueue_handler_with_options(
1627 "echo-workflow",
1628 TriggerKind::Cron {
1629 schedule: "0 * * * * *".to_string(),
1630 },
1631 json!({}),
1632 EnqueueOptions::default(),
1633 )
1634 .await
1635 .unwrap()
1636 .into_run();
1637
1638 assert!(run.created_by.is_none());
1639 }
1640
1641 #[tokio::test]
1642 async fn run_handler_leaves_the_run_unattributed() {
1643 let mut engine = create_test_engine();
1644 engine.register(EchoWorkflow).unwrap();
1645
1646 let run = engine
1647 .run_handler("echo-workflow", TriggerKind::Manual, json!({}))
1648 .await
1649 .unwrap();
1650
1651 assert!(run.created_by.is_none());
1652 }
1653
1654 #[tokio::test]
1655 async fn engine_register_boxed() {
1656 let mut engine = create_test_engine();
1657 let handler: Box<dyn WorkflowHandler> = Box::new(EchoWorkflow);
1658 let result = engine.register_boxed(handler);
1659 assert!(result.is_ok());
1660 assert_eq!(engine.handler_names().len(), 1);
1661 }
1662
1663 #[tokio::test]
1664 async fn engine_store_and_provider_accessors() {
1665 let store = Arc::new(InMemoryStore::new());
1666 let inner = ClaudeCodeProvider::new();
1667 let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
1668 inner,
1669 "/tmp/ironflow-fixtures",
1670 ));
1671 let engine = Engine::new(store.clone(), provider.clone());
1672
1673 let _ = engine.store();
1675 let _ = engine.provider();
1676 }
1677
1678 use crate::operation::Operation;
1683 use ironflow_store::models::StepKind;
1684 use std::future::Future;
1685 use std::pin::Pin;
1686
1687 struct FakeGitlabOp {
1688 project_id: u64,
1689 title: String,
1690 }
1691
1692 impl Operation for FakeGitlabOp {
1693 fn kind(&self) -> &str {
1694 "gitlab"
1695 }
1696
1697 fn execute(&self) -> Pin<Box<dyn Future<Output = Result<Value, EngineError>> + Send + '_>> {
1698 Box::pin(async move {
1699 Ok(json!({
1700 "issue_id": 42,
1701 "project_id": self.project_id,
1702 "title": self.title,
1703 }))
1704 })
1705 }
1706
1707 fn input(&self) -> Option<Value> {
1708 Some(json!({
1709 "project_id": self.project_id,
1710 "title": self.title,
1711 }))
1712 }
1713 }
1714
1715 struct FailingOp;
1716
1717 impl Operation for FailingOp {
1718 fn kind(&self) -> &str {
1719 "broken-service"
1720 }
1721
1722 fn execute(&self) -> Pin<Box<dyn Future<Output = Result<Value, EngineError>> + Send + '_>> {
1723 Box::pin(async move { Err(EngineError::StepConfig("service unavailable".to_string())) })
1724 }
1725 }
1726
1727 struct OperationWorkflow;
1728
1729 impl WorkflowHandler for OperationWorkflow {
1730 fn name(&self) -> &str {
1731 "operation-workflow"
1732 }
1733
1734 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
1735 Box::pin(async move {
1736 let op = FakeGitlabOp {
1737 project_id: 123,
1738 title: "Bug report".to_string(),
1739 };
1740 ctx.operation("create-issue", &op).await?;
1741 Ok(())
1742 })
1743 }
1744 }
1745
1746 struct FailingOperationWorkflow;
1747
1748 impl WorkflowHandler for FailingOperationWorkflow {
1749 fn name(&self) -> &str {
1750 "failing-operation-workflow"
1751 }
1752
1753 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
1754 Box::pin(async move {
1755 ctx.operation("broken-call", &FailingOp).await?;
1756 Ok(())
1757 })
1758 }
1759 }
1760
1761 struct MixedWorkflow;
1762
1763 impl WorkflowHandler for MixedWorkflow {
1764 fn name(&self) -> &str {
1765 "mixed-workflow"
1766 }
1767
1768 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
1769 Box::pin(async move {
1770 ctx.shell("build", ShellConfig::new("echo built")).await?;
1771 let op = FakeGitlabOp {
1772 project_id: 456,
1773 title: "Deploy done".to_string(),
1774 };
1775 let result = ctx.operation("notify-gitlab", &op).await?;
1776 assert_eq!(result.output["issue_id"], 42);
1777 Ok(())
1778 })
1779 }
1780 }
1781
1782 #[tokio::test]
1783 async fn operation_step_happy_path() {
1784 let mut engine = create_test_engine();
1785 engine.register(OperationWorkflow).unwrap();
1786
1787 let run = engine
1788 .run_handler("operation-workflow", TriggerKind::Manual, json!({}))
1789 .await
1790 .unwrap();
1791
1792 assert_eq!(run.status.state, RunStatus::Completed);
1793
1794 let steps = engine.store().list_steps(run.id).await.unwrap();
1795
1796 assert_eq!(steps.len(), 1);
1797 assert_eq!(steps[0].name, "create-issue");
1798 assert_eq!(steps[0].kind, StepKind::Custom("gitlab".to_string()));
1799 assert_eq!(
1800 steps[0].status.state,
1801 ironflow_store::models::StepStatus::Completed
1802 );
1803
1804 let output = steps[0].output.as_ref().unwrap();
1805 assert_eq!(output["issue_id"], 42);
1806 assert_eq!(output["project_id"], 123);
1807
1808 let input = steps[0].input.as_ref().unwrap();
1809 assert_eq!(input["project_id"], 123);
1810 assert_eq!(input["title"], "Bug report");
1811 }
1812
1813 #[tokio::test]
1814 async fn operation_step_failure_marks_run_failed() {
1815 let mut engine = create_test_engine();
1816 engine.register(FailingOperationWorkflow).unwrap();
1817
1818 let result = engine
1819 .run_handler("failing-operation-workflow", TriggerKind::Manual, json!({}))
1820 .await;
1821
1822 assert!(result.is_err());
1823 }
1824
1825 #[tokio::test]
1826 async fn operation_mixed_with_shell_steps() {
1827 let mut engine = create_test_engine();
1828 engine.register(MixedWorkflow).unwrap();
1829
1830 let run = engine
1831 .run_handler("mixed-workflow", TriggerKind::Manual, json!({}))
1832 .await
1833 .unwrap();
1834
1835 assert_eq!(run.status.state, RunStatus::Completed);
1836
1837 let steps = engine.store().list_steps(run.id).await.unwrap();
1838
1839 assert_eq!(steps.len(), 2);
1840 assert_eq!(steps[0].kind, StepKind::Shell);
1841 assert_eq!(steps[1].kind, StepKind::Custom("gitlab".to_string()));
1842 assert_eq!(steps[0].position, 0);
1843 assert_eq!(steps[1].position, 1);
1844 }
1845
1846 use crate::config::ApprovalConfig;
1851
1852 struct SingleApprovalWorkflow;
1853
1854 impl WorkflowHandler for SingleApprovalWorkflow {
1855 fn name(&self) -> &str {
1856 "single-approval"
1857 }
1858
1859 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
1860 Box::pin(async move {
1861 ctx.shell("build", ShellConfig::new("echo built")).await?;
1862 ctx.approval("gate", ApprovalConfig::new("OK?")).await?;
1863 ctx.shell("deploy", ShellConfig::new("echo deployed"))
1864 .await?;
1865 Ok(())
1866 })
1867 }
1868 }
1869
1870 struct DoubleApprovalWorkflow;
1871
1872 impl WorkflowHandler for DoubleApprovalWorkflow {
1873 fn name(&self) -> &str {
1874 "double-approval"
1875 }
1876
1877 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
1878 Box::pin(async move {
1879 ctx.shell("build", ShellConfig::new("echo built")).await?;
1880 ctx.approval("staging-gate", ApprovalConfig::new("Deploy staging?"))
1881 .await?;
1882 ctx.shell("deploy-staging", ShellConfig::new("echo staging"))
1883 .await?;
1884 ctx.approval("prod-gate", ApprovalConfig::new("Deploy prod?"))
1885 .await?;
1886 ctx.shell("deploy-prod", ShellConfig::new("echo prod"))
1887 .await?;
1888 Ok(())
1889 })
1890 }
1891 }
1892
1893 #[tokio::test]
1894 async fn approval_pauses_run() {
1895 let mut engine = create_test_engine();
1896 engine.register(SingleApprovalWorkflow).unwrap();
1897
1898 let run = engine
1899 .run_handler("single-approval", TriggerKind::Manual, json!({}))
1900 .await
1901 .unwrap();
1902
1903 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
1904
1905 let steps = engine.store().list_steps(run.id).await.unwrap();
1906 assert_eq!(steps.len(), 2); assert_eq!(steps[0].kind, StepKind::Shell);
1908 assert_eq!(steps[0].status.state, StepStatus::Completed);
1909 assert_eq!(steps[1].kind, StepKind::Approval);
1910 assert_eq!(steps[1].status.state, StepStatus::AwaitingApproval);
1911 }
1912
1913 #[tokio::test]
1914 async fn approval_resume_completes_run() {
1915 let mut engine = create_test_engine();
1916 engine.register(SingleApprovalWorkflow).unwrap();
1917
1918 let run = engine
1920 .run_handler("single-approval", TriggerKind::Manual, json!({}))
1921 .await
1922 .unwrap();
1923 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
1924
1925 engine
1927 .store()
1928 .update_run_status(run.id, RunStatus::Running)
1929 .await
1930 .unwrap();
1931
1932 let resumed = engine.resume_run(run.id).await.unwrap();
1934 assert_eq!(resumed.status.state, RunStatus::Completed);
1935
1936 let steps = engine.store().list_steps(run.id).await.unwrap();
1937 assert_eq!(steps.len(), 3); assert_eq!(steps[0].name, "build");
1939 assert_eq!(steps[0].status.state, StepStatus::Completed);
1940 assert_eq!(steps[1].name, "gate");
1941 assert_eq!(steps[1].kind, StepKind::Approval);
1942 assert_eq!(steps[1].status.state, StepStatus::Completed);
1943 assert_eq!(steps[2].name, "deploy");
1944 assert_eq!(steps[2].status.state, StepStatus::Completed);
1945 }
1946
1947 #[tokio::test]
1948 async fn double_approval_two_resumes() {
1949 let mut engine = create_test_engine();
1950 engine.register(DoubleApprovalWorkflow).unwrap();
1951
1952 let run = engine
1954 .run_handler("double-approval", TriggerKind::Manual, json!({}))
1955 .await
1956 .unwrap();
1957 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
1958
1959 let steps = engine.store().list_steps(run.id).await.unwrap();
1960 assert_eq!(steps.len(), 2); engine
1964 .store()
1965 .update_run_status(run.id, RunStatus::Running)
1966 .await
1967 .unwrap();
1968
1969 let resumed = engine.resume_run(run.id).await.unwrap();
1970 assert_eq!(resumed.status.state, RunStatus::AwaitingApproval);
1971
1972 let steps = engine.store().list_steps(run.id).await.unwrap();
1973 assert_eq!(steps.len(), 4); engine
1977 .store()
1978 .update_run_status(run.id, RunStatus::Running)
1979 .await
1980 .unwrap();
1981
1982 let final_run = engine.resume_run(run.id).await.unwrap();
1983 assert_eq!(final_run.status.state, RunStatus::Completed);
1984
1985 let steps = engine.store().list_steps(run.id).await.unwrap();
1986 assert_eq!(steps.len(), 5);
1987 assert_eq!(steps[0].name, "build");
1988 assert_eq!(steps[1].name, "staging-gate");
1989 assert_eq!(steps[2].name, "deploy-staging");
1990 assert_eq!(steps[3].name, "prod-gate");
1991 assert_eq!(steps[4].name, "deploy-prod");
1992
1993 for step in &steps {
1994 assert_eq!(step.status.state, StepStatus::Completed);
1995 }
1996 }
1997
1998 use ironflow_store::models::{NewStep, StepUpdate};
2003
2004 async fn create_step_with_status(
2005 store: &Arc<dyn Store>,
2006 run_id: Uuid,
2007 name: &str,
2008 position: u32,
2009 status: StepStatus,
2010 ) -> ironflow_store::models::Step {
2011 let step = store
2012 .create_step(NewStep {
2013 run_id,
2014 name: name.to_string(),
2015 kind: StepKind::Shell,
2016 position,
2017 input: None,
2018 })
2019 .await
2020 .unwrap();
2021
2022 match status {
2023 StepStatus::Pending => {}
2024 StepStatus::Running => {
2025 store
2026 .update_step(
2027 step.id,
2028 StepUpdate {
2029 status: Some(StepStatus::Running),
2030 ..StepUpdate::default()
2031 },
2032 )
2033 .await
2034 .unwrap();
2035 }
2036 StepStatus::Completed => {
2037 store
2038 .update_step(
2039 step.id,
2040 StepUpdate {
2041 status: Some(StepStatus::Running),
2042 ..StepUpdate::default()
2043 },
2044 )
2045 .await
2046 .unwrap();
2047 store
2048 .update_step(
2049 step.id,
2050 StepUpdate {
2051 status: Some(StepStatus::Completed),
2052 ..StepUpdate::default()
2053 },
2054 )
2055 .await
2056 .unwrap();
2057 }
2058 StepStatus::AwaitingApproval => {
2059 store
2060 .update_step(
2061 step.id,
2062 StepUpdate {
2063 status: Some(StepStatus::Running),
2064 ..StepUpdate::default()
2065 },
2066 )
2067 .await
2068 .unwrap();
2069 store
2070 .update_step(
2071 step.id,
2072 StepUpdate {
2073 status: Some(StepStatus::AwaitingApproval),
2074 ..StepUpdate::default()
2075 },
2076 )
2077 .await
2078 .unwrap();
2079 }
2080 _ => panic!("unsupported status for test helper: {status}"),
2081 }
2082
2083 store.get_step(step.id).await.unwrap().unwrap()
2084 }
2085
2086 #[tokio::test]
2087 async fn fail_orphaned_steps_marks_running_as_failed() {
2088 let engine = create_test_engine();
2089 let run = engine
2090 .store()
2091 .create_run(NewRun {
2092 created_by: None,
2093 workflow_name: "test".to_string(),
2094 trigger: TriggerKind::Manual,
2095 payload: json!({}),
2096 max_retries: 0,
2097 handler_version: None,
2098 labels: HashMap::new(),
2099 scheduled_at: None,
2100 idempotency_key: None,
2101 max_cost_usd: None,
2102 })
2103 .await
2104 .unwrap()
2105 .into_run();
2106
2107 let step = create_step_with_status(
2108 engine.store(),
2109 run.id,
2110 "running-step",
2111 0,
2112 StepStatus::Running,
2113 )
2114 .await;
2115
2116 engine
2117 .fail_orphaned_steps(run.id, "parent run timed out")
2118 .await
2119 .unwrap();
2120
2121 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
2122 assert_eq!(updated.status.state, StepStatus::Failed);
2123 assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
2124 assert!(updated.completed_at.is_some());
2125 }
2126
2127 #[tokio::test]
2128 async fn fail_orphaned_steps_marks_pending_as_skipped() {
2129 let engine = create_test_engine();
2130 let run = engine
2131 .store()
2132 .create_run(NewRun {
2133 created_by: None,
2134 workflow_name: "test".to_string(),
2135 trigger: TriggerKind::Manual,
2136 payload: json!({}),
2137 max_retries: 0,
2138 handler_version: None,
2139 labels: HashMap::new(),
2140 scheduled_at: None,
2141 idempotency_key: None,
2142 max_cost_usd: None,
2143 })
2144 .await
2145 .unwrap()
2146 .into_run();
2147
2148 let step = create_step_with_status(
2149 engine.store(),
2150 run.id,
2151 "pending-step",
2152 0,
2153 StepStatus::Pending,
2154 )
2155 .await;
2156
2157 engine
2158 .fail_orphaned_steps(run.id, "parent run timed out")
2159 .await
2160 .unwrap();
2161
2162 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
2163 assert_eq!(updated.status.state, StepStatus::Skipped);
2164 assert!(updated.error.is_none());
2165 assert!(updated.completed_at.is_some());
2166 }
2167
2168 #[tokio::test]
2169 async fn fail_orphaned_steps_marks_awaiting_approval_as_failed() {
2170 let engine = create_test_engine();
2171 let run = engine
2172 .store()
2173 .create_run(NewRun {
2174 created_by: None,
2175 workflow_name: "test".to_string(),
2176 trigger: TriggerKind::Manual,
2177 payload: json!({}),
2178 max_retries: 0,
2179 handler_version: None,
2180 labels: HashMap::new(),
2181 scheduled_at: None,
2182 idempotency_key: None,
2183 max_cost_usd: None,
2184 })
2185 .await
2186 .unwrap()
2187 .into_run();
2188
2189 let step = create_step_with_status(
2190 engine.store(),
2191 run.id,
2192 "approval-step",
2193 0,
2194 StepStatus::AwaitingApproval,
2195 )
2196 .await;
2197
2198 engine
2199 .fail_orphaned_steps(run.id, "parent run timed out")
2200 .await
2201 .unwrap();
2202
2203 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
2204 assert_eq!(updated.status.state, StepStatus::Failed);
2205 assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
2206 assert!(updated.completed_at.is_some());
2207 }
2208
2209 #[tokio::test]
2210 async fn fail_orphaned_steps_skips_terminal_steps() {
2211 let engine = create_test_engine();
2212 let run = engine
2213 .store()
2214 .create_run(NewRun {
2215 created_by: None,
2216 workflow_name: "test".to_string(),
2217 trigger: TriggerKind::Manual,
2218 payload: json!({}),
2219 max_retries: 0,
2220 handler_version: None,
2221 labels: HashMap::new(),
2222 scheduled_at: None,
2223 idempotency_key: None,
2224 max_cost_usd: None,
2225 })
2226 .await
2227 .unwrap()
2228 .into_run();
2229
2230 let completed_step =
2231 create_step_with_status(engine.store(), run.id, "done", 0, StepStatus::Completed).await;
2232 let running_step =
2233 create_step_with_status(engine.store(), run.id, "in-flight", 1, StepStatus::Running)
2234 .await;
2235
2236 engine
2237 .fail_orphaned_steps(run.id, "parent run timed out")
2238 .await
2239 .unwrap();
2240
2241 let completed = engine
2242 .store()
2243 .get_step(completed_step.id)
2244 .await
2245 .unwrap()
2246 .unwrap();
2247 assert_eq!(completed.status.state, StepStatus::Completed);
2248
2249 let failed = engine
2250 .store()
2251 .get_step(running_step.id)
2252 .await
2253 .unwrap()
2254 .unwrap();
2255 assert_eq!(failed.status.state, StepStatus::Failed);
2256 }
2257
2258 #[tokio::test]
2259 async fn fail_orphaned_steps_mixed_states() {
2260 let engine = create_test_engine();
2261 let run = engine
2262 .store()
2263 .create_run(NewRun {
2264 created_by: None,
2265 workflow_name: "test".to_string(),
2266 trigger: TriggerKind::Manual,
2267 payload: json!({}),
2268 max_retries: 0,
2269 handler_version: None,
2270 labels: HashMap::new(),
2271 scheduled_at: None,
2272 idempotency_key: None,
2273 max_cost_usd: None,
2274 })
2275 .await
2276 .unwrap()
2277 .into_run();
2278
2279 let s_completed =
2280 create_step_with_status(engine.store(), run.id, "step-1", 0, StepStatus::Completed)
2281 .await;
2282 let s_running =
2283 create_step_with_status(engine.store(), run.id, "step-2", 1, StepStatus::Running).await;
2284 let s_pending =
2285 create_step_with_status(engine.store(), run.id, "step-3", 2, StepStatus::Pending).await;
2286
2287 engine.fail_orphaned_steps(run.id, "timeout").await.unwrap();
2288
2289 let r_completed = engine
2290 .store()
2291 .get_step(s_completed.id)
2292 .await
2293 .unwrap()
2294 .unwrap();
2295 assert_eq!(r_completed.status.state, StepStatus::Completed);
2296
2297 let r_running = engine
2298 .store()
2299 .get_step(s_running.id)
2300 .await
2301 .unwrap()
2302 .unwrap();
2303 assert_eq!(r_running.status.state, StepStatus::Failed);
2304 assert_eq!(r_running.error.as_deref(), Some("timeout"));
2305
2306 let r_pending = engine
2307 .store()
2308 .get_step(s_pending.id)
2309 .await
2310 .unwrap()
2311 .unwrap();
2312 assert_eq!(r_pending.status.state, StepStatus::Skipped);
2313 assert!(r_pending.error.is_none());
2314 }
2315
2316 #[tokio::test]
2317 async fn fail_orphaned_steps_no_steps_is_noop() {
2318 let engine = create_test_engine();
2319 let run = engine
2320 .store()
2321 .create_run(NewRun {
2322 created_by: None,
2323 workflow_name: "test".to_string(),
2324 trigger: TriggerKind::Manual,
2325 payload: json!({}),
2326 max_retries: 0,
2327 handler_version: None,
2328 labels: HashMap::new(),
2329 scheduled_at: None,
2330 idempotency_key: None,
2331 max_cost_usd: None,
2332 })
2333 .await
2334 .unwrap()
2335 .into_run();
2336
2337 let result = engine.fail_orphaned_steps(run.id, "timeout").await;
2338 assert!(result.is_ok());
2339 }
2340
2341 #[tokio::test]
2342 async fn fail_orphaned_steps_preserves_existing_error() {
2343 let engine = create_test_engine();
2344 let run = engine
2345 .store()
2346 .create_run(NewRun {
2347 created_by: None,
2348 workflow_name: "test".to_string(),
2349 trigger: TriggerKind::Manual,
2350 payload: json!({}),
2351 max_retries: 0,
2352 handler_version: None,
2353 labels: HashMap::new(),
2354 scheduled_at: None,
2355 idempotency_key: None,
2356 max_cost_usd: None,
2357 })
2358 .await
2359 .unwrap()
2360 .into_run();
2361
2362 let step_with_error = create_step_with_status(
2363 engine.store(),
2364 run.id,
2365 "already-errored",
2366 0,
2367 StepStatus::Running,
2368 )
2369 .await;
2370
2371 engine
2372 .store()
2373 .update_step(
2374 step_with_error.id,
2375 StepUpdate {
2376 error: Some("real error from provider".to_string()),
2377 ..StepUpdate::default()
2378 },
2379 )
2380 .await
2381 .unwrap();
2382
2383 let step_no_error = create_step_with_status(
2384 engine.store(),
2385 run.id,
2386 "no-error-yet",
2387 1,
2388 StepStatus::Running,
2389 )
2390 .await;
2391
2392 engine
2393 .fail_orphaned_steps(run.id, "parent run failed")
2394 .await
2395 .unwrap();
2396
2397 let updated_with = engine
2398 .store()
2399 .get_step(step_with_error.id)
2400 .await
2401 .unwrap()
2402 .unwrap();
2403 assert_eq!(updated_with.status.state, StepStatus::Failed);
2404 assert_eq!(
2405 updated_with.error.as_deref(),
2406 Some("real error from provider"),
2407 );
2408
2409 let updated_without = engine
2410 .store()
2411 .get_step(step_no_error.id)
2412 .await
2413 .unwrap()
2414 .unwrap();
2415 assert_eq!(updated_without.status.state, StepStatus::Failed);
2416 assert_eq!(updated_without.error.as_deref(), Some("parent run failed"),);
2417 }
2418}