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::artifact::ArtifactSink;
35use crate::budget::{BudgetConfig, month_start};
36use crate::context::WorkflowContext;
37use crate::error::EngineError;
38use crate::handler::{WorkflowHandler, WorkflowInfo};
39use crate::log_sender::LogSender;
40use crate::notify::{Event, EventPublisher, EventSubscriber};
41use crate::retry_policy::{backoff_for_retry, is_run_retryable};
42use crate::schedule::CronSchedule;
43
44#[derive(Debug, Clone, Default)]
63pub struct EnqueueOptions {
64 pub max_retries: u32,
66 pub labels: HashMap<String, String>,
68 pub scheduled_at: Option<DateTime<Utc>>,
71 pub max_cost_usd: Option<Decimal>,
75 pub created_by: Option<RunActor>,
78 pub idempotency_key: Option<String>,
84}
85
86pub struct Engine {
127 store: Arc<dyn Store>,
128 provider: Arc<dyn AgentProvider>,
129 handlers: HashMap<String, Arc<dyn WorkflowHandler>>,
130 event_publisher: EventPublisher,
131 log_sender: Option<LogSender>,
132 budget: BudgetConfig,
133 artifact_sink: Option<Arc<dyn ArtifactSink>>,
134}
135
136fn validate_category(handler_name: &str, category: &str) -> Result<(), EngineError> {
146 let reject = |reason: &str| {
147 Err(EngineError::InvalidWorkflow(format!(
148 "handler '{handler_name}' has invalid category '{category}': {reason}"
149 )))
150 };
151
152 if category.is_empty() {
153 return reject("empty category");
154 }
155 if category.starts_with('/') {
156 return reject("leading '/'");
157 }
158 if category.ends_with('/') {
159 return reject("trailing '/'");
160 }
161 for segment in category.split('/') {
162 if segment.is_empty() {
163 return reject("empty segment (double '/')");
164 }
165 if segment.trim().is_empty() {
166 return reject("whitespace-only segment");
167 }
168 }
169 Ok(())
170}
171
172impl Engine {
173 pub fn new(store: Arc<dyn Store>, provider: Arc<dyn AgentProvider>) -> Self {
189 Self {
190 store,
191 provider,
192 handlers: HashMap::new(),
193 event_publisher: EventPublisher::new(),
194 log_sender: None,
195 budget: BudgetConfig::new(),
196 artifact_sink: None,
197 }
198 }
199
200 pub fn with_budget_config(mut self, budget: BudgetConfig) -> Self {
221 self.budget = budget;
222 self
223 }
224
225 pub fn budget_config(&self) -> &BudgetConfig {
227 &self.budget
228 }
229
230 pub fn set_log_sender(&mut self, sender: LogSender) {
236 self.log_sender = Some(sender);
237 }
238
239 pub fn set_artifact_sink(&mut self, sink: Arc<dyn ArtifactSink>) {
257 self.artifact_sink = Some(sink);
258 }
259
260 pub fn artifact_sink(&self) -> Option<&Arc<dyn ArtifactSink>> {
262 self.artifact_sink.as_ref()
263 }
264
265 pub fn store(&self) -> &Arc<dyn Store> {
267 &self.store
268 }
269
270 pub fn provider(&self) -> &Arc<dyn AgentProvider> {
272 &self.provider
273 }
274
275 fn build_context(&self, run: &Run) -> WorkflowContext {
284 let handlers = self.handlers.clone();
285 let resolver: crate::context::HandlerResolver =
286 Arc::new(move |name: &str| handlers.get(name).cloned());
287 let mut ctx = WorkflowContext::with_handler_resolver(
288 run.id,
289 self.store.clone(),
290 self.provider.clone(),
291 resolver,
292 );
293 ctx.carry_over_run_totals(run.retry_count + 1, run.cost_usd, run.duration_ms);
294 ctx.set_max_cost_usd(run.max_cost_usd);
295 if let Some(ref sender) = self.log_sender {
296 ctx.set_log_sender(sender.clone());
297 }
298 if let Some(ref sink) = self.artifact_sink {
299 ctx.set_artifact_sink(sink.clone());
300 }
301 ctx
302 }
303
304 async fn check_monthly_quota(&self, workflow_name: &str) -> Result<(), EngineError> {
315 let Some(limit) = self.budget.monthly_cost_limit_usd else {
316 return Ok(());
317 };
318
319 let stats = self
320 .store
321 .get_stats(RunFilter {
322 created_after: Some(month_start(Utc::now())),
323 ..RunFilter::default()
324 })
325 .await?;
326
327 if stats.total_cost_usd < limit {
328 return Ok(());
329 }
330
331 warn!(
332 workflow = %workflow_name,
333 limit_usd = %limit,
334 spent_usd = %stats.total_cost_usd,
335 "monthly cost quota exhausted, refusing new run"
336 );
337
338 #[cfg(feature = "prometheus")]
339 counter!(
340 RUN_BUDGET_EXCEEDED_TOTAL,
341 "workflow" => workflow_name.to_string(),
342 "scope" => "monthly",
343 )
344 .increment(1);
345
346 Err(EngineError::MonthlyBudgetExceeded {
347 limit_usd: limit,
348 spent_usd: stats.total_cost_usd,
349 })
350 }
351
352 pub fn register(&mut self, handler: impl WorkflowHandler + 'static) -> Result<(), EngineError> {
396 let name = handler.name().to_string();
397 if self.handlers.contains_key(&name) {
398 return Err(EngineError::InvalidWorkflow(format!(
399 "handler '{}' already registered",
400 name
401 )));
402 }
403 if let Some(category) = handler.category() {
404 validate_category(&name, category)?;
405 }
406 self.handlers.insert(name, Arc::new(handler));
407 Ok(())
408 }
409
410 pub fn register_boxed(&mut self, handler: Box<dyn WorkflowHandler>) -> Result<(), EngineError> {
417 let name = handler.name().to_string();
418 if self.handlers.contains_key(&name) {
419 return Err(EngineError::InvalidWorkflow(format!(
420 "handler '{}' already registered",
421 name
422 )));
423 }
424 if let Some(category) = handler.category() {
425 validate_category(&name, category)?;
426 }
427 self.handlers.insert(name, Arc::from(handler));
428 Ok(())
429 }
430
431 pub fn get_handler(&self, name: &str) -> Option<&Arc<dyn WorkflowHandler>> {
433 self.handlers.get(name)
434 }
435
436 pub fn handler_names(&self) -> Vec<&str> {
438 self.handlers.keys().map(|s| s.as_str()).collect()
439 }
440
441 pub fn handler_info(&self, name: &str) -> Option<WorkflowInfo> {
443 self.handlers.get(name).map(|h| h.describe())
444 }
445
446 pub fn scheduled_handlers(&self) -> Vec<(&str, &CronSchedule)> {
470 self.handlers
471 .iter()
472 .filter_map(|(name, handler)| handler.schedule().map(|sched| (name.as_str(), sched)))
473 .collect()
474 }
475
476 pub fn subscribe(
501 &mut self,
502 subscriber: impl EventSubscriber + 'static,
503 event_types: &[&'static str],
504 ) {
505 self.event_publisher.subscribe(subscriber, event_types);
506 }
507
508 pub fn event_publisher(&self) -> &EventPublisher {
513 &self.event_publisher
514 }
515
516 #[tracing::instrument(name = "engine.run_handler", skip_all, fields(workflow = %handler_name))]
546 pub async fn run_handler(
547 &self,
548 handler_name: &str,
549 trigger: TriggerKind,
550 payload: Value,
551 ) -> Result<Run, EngineError> {
552 let handler = self
553 .handlers
554 .get(handler_name)
555 .ok_or_else(|| {
556 EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
557 })?
558 .clone();
559
560 self.check_monthly_quota(handler_name).await?;
561
562 let handler_version = handler.version().map(str::to_string);
563 let max_cost_usd = self
564 .budget
565 .resolve_run_cap(None, handler.default_max_cost_usd());
566 let run = self
567 .store
568 .create_run(NewRun {
569 created_by: None,
570 workflow_name: handler_name.to_string(),
571 trigger,
572 payload,
573 max_retries: 0,
574 handler_version,
575 labels: handler.default_labels(),
576 scheduled_at: None,
577 idempotency_key: None,
578 max_cost_usd,
579 })
580 .await?
581 .into_run();
582
583 let run_id = run.id;
584 info!(run_id = %run_id, handler_version = run.handler_version.as_deref().unwrap_or(""), "run created");
585
586 self.store
587 .update_run_status(run_id, RunStatus::Running)
588 .await?;
589
590 #[cfg(feature = "prometheus")]
591 gauge!(RUNS_ACTIVE, "workflow" => handler_name.to_string()).increment(1.0);
592
593 let run_start = Instant::now();
594 let mut ctx = self.build_context(&run);
595
596 let result = handler.execute(&mut ctx).await;
597 self.finalize_run(run_id, handler_name, result, &ctx, run_start, run.labels)
598 .await
599 }
600
601 #[tracing::instrument(name = "engine.enqueue_handler", skip_all, fields(workflow = %handler_name))]
612 pub async fn enqueue_handler(
613 &self,
614 handler_name: &str,
615 trigger: TriggerKind,
616 payload: Value,
617 max_retries: u32,
618 ) -> Result<Run, EngineError> {
619 self.enqueue_handler_with_options(
620 handler_name,
621 trigger,
622 payload,
623 EnqueueOptions {
624 max_retries,
625 ..Default::default()
626 },
627 )
628 .await
629 .map(RunCreation::into_run)
630 }
631
632 #[tracing::instrument(name = "engine.enqueue_handler_with_options", skip_all, fields(workflow = %handler_name))]
677 pub async fn enqueue_handler_with_options(
678 &self,
679 handler_name: &str,
680 trigger: TriggerKind,
681 payload: Value,
682 options: EnqueueOptions,
683 ) -> Result<RunCreation, EngineError> {
684 let EnqueueOptions {
685 max_retries,
686 labels,
687 scheduled_at,
688 max_cost_usd,
689 created_by,
690 idempotency_key,
691 } = options;
692
693 let handler = self.handlers.get(handler_name).ok_or_else(|| {
694 EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
695 })?;
696
697 self.check_monthly_quota(handler_name).await?;
698
699 let handler_version = handler.version().map(str::to_string);
700 let mut merged_labels = handler.default_labels();
701 merged_labels.extend(labels);
702 let resolved_cap = self
703 .budget
704 .resolve_run_cap(max_cost_usd, handler.default_max_cost_usd());
705
706 let creation = self
707 .store
708 .create_run(NewRun {
709 workflow_name: handler_name.to_string(),
710 trigger,
711 payload,
712 max_retries,
713 handler_version,
714 labels: merged_labels,
715 scheduled_at,
716 created_by,
717 idempotency_key,
718 max_cost_usd: resolved_cap,
719 })
720 .await?;
721
722 match &creation {
723 RunCreation::Created(run) => info!(
724 run_id = %run.id,
725 workflow = %handler_name,
726 max_cost_usd = ?resolved_cap,
727 "handler run enqueued"
728 ),
729 RunCreation::Existing(run) => info!(
730 run_id = %run.id,
731 workflow = %handler_name,
732 "idempotent replay, nothing enqueued"
733 ),
734 }
735
736 Ok(creation)
737 }
738
739 #[tracing::instrument(name = "engine.execute_handler_run", skip_all, fields(run_id = %run_id))]
748 pub async fn execute_handler_run(&self, run_id: Uuid) -> Result<Run, EngineError> {
749 let run = self
750 .store
751 .get_run(run_id)
752 .await?
753 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
754
755 let handler = self
756 .handlers
757 .get(&run.workflow_name)
758 .ok_or_else(|| {
759 EngineError::InvalidWorkflow(format!(
760 "no handler registered: {}",
761 run.workflow_name
762 ))
763 })?
764 .clone();
765
766 #[cfg(feature = "prometheus")]
767 gauge!(RUNS_ACTIVE, "workflow" => run.workflow_name.clone()).increment(1.0);
768
769 let run_start = Instant::now();
770 let mut ctx = self.build_context(&run);
771
772 if run.retry_count > 0 {
775 ctx.load_replay_steps().await?;
776 }
777
778 let result = handler.execute(&mut ctx).await;
779 self.finalize_run(
780 run_id,
781 &run.workflow_name,
782 result,
783 &ctx,
784 run_start,
785 run.labels,
786 )
787 .await
788 }
789
790 #[tracing::instrument(name = "engine.execute_run", skip_all, fields(run_id = %run_id))]
798 pub async fn execute_run(&self, run_id: Uuid) -> Result<Run, EngineError> {
799 self.execute_handler_run(run_id).await
800 }
801
802 #[tracing::instrument(name = "engine.resume_run", skip_all, fields(run_id = %run_id))]
816 pub async fn resume_run(&self, run_id: Uuid) -> Result<Run, EngineError> {
817 let run = self
818 .store
819 .get_run(run_id)
820 .await?
821 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
822
823 let handler = self
824 .handlers
825 .get(&run.workflow_name)
826 .ok_or_else(|| {
827 EngineError::InvalidWorkflow(format!(
828 "no handler registered: {}",
829 run.workflow_name
830 ))
831 })?
832 .clone();
833
834 info!(run_id = %run_id, workflow = %run.workflow_name, "resuming run after approval");
835
836 let run_start = Instant::now();
837 let mut ctx = self.build_context(&run);
838 ctx.load_replay_steps().await?;
839
840 let result = handler.execute(&mut ctx).await;
841 self.finalize_run(
842 run_id,
843 &run.workflow_name,
844 result,
845 &ctx,
846 run_start,
847 run.labels,
848 )
849 .await
850 }
851
852 pub async fn fail_or_schedule_retry(
897 &self,
898 run_id: Uuid,
899 error: &str,
900 retryable: bool,
901 cost_usd: Option<Decimal>,
902 duration_ms: Option<u64>,
903 ) -> Result<RunStatus, EngineError> {
904 let run = self
905 .store
906 .get_run(run_id)
907 .await?
908 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
909
910 let has_attempts_left = run.retry_count < run.max_retries;
911 let update = if retryable && has_attempts_left {
912 let backoff = backoff_for_retry(run.retry_count);
913 let scheduled_at = Utc::now() + TimeDelta::milliseconds(backoff.as_millis() as i64);
914
915 info!(
916 run_id = %run_id,
917 workflow = %run.workflow_name,
918 attempt = run.retry_count + 1,
919 max_retries = run.max_retries,
920 backoff_secs = backoff.as_secs(),
921 scheduled_at = %scheduled_at,
922 "run failed, scheduling retry"
923 );
924
925 RunUpdate {
926 status: Some(RunStatus::Retrying),
927 error: Some(error.to_string()),
928 increment_retry: true,
929 cost_usd,
930 duration_ms,
931 scheduled_at: Some(scheduled_at),
932 ..RunUpdate::default()
933 }
934 } else {
935 RunUpdate {
936 status: Some(RunStatus::Failed),
937 error: Some(error.to_string()),
938 cost_usd,
939 duration_ms,
940 completed_at: Some(Utc::now()),
941 ..RunUpdate::default()
942 }
943 };
944
945 let status = update.status.unwrap_or(RunStatus::Failed);
946 self.store.update_run(run_id, update).await?;
947 self.fail_orphaned_steps(run_id, error).await?;
948
949 Ok(status)
950 }
951
952 pub async fn fail_orphaned_steps(
966 &self,
967 run_id: Uuid,
968 error_message: &str,
969 ) -> Result<(), EngineError> {
970 let steps = self.store.list_steps(run_id).await?;
971 let now = Utc::now();
972
973 for step in steps {
974 if step.status.state.is_terminal() {
975 continue;
976 }
977
978 let (target_status, error) = match step.status.state {
979 StepStatus::Running | StepStatus::AwaitingApproval => {
980 let err = if step.error.is_some() {
981 None
982 } else {
983 Some(error_message.to_string())
984 };
985 (StepStatus::Failed, err)
986 }
987 StepStatus::Pending => (StepStatus::Skipped, None),
988 _ => continue,
989 };
990
991 if let Err(e) = self
992 .store
993 .update_step(
994 step.id,
995 StepUpdate {
996 status: Some(target_status),
997 error,
998 completed_at: Some(now),
999 ..StepUpdate::default()
1000 },
1001 )
1002 .await
1003 {
1004 warn!(
1005 run_id = %run_id,
1006 step_id = %step.id,
1007 step_name = %step.name,
1008 error = %e,
1009 "failed to cleanup orphaned step"
1010 );
1011 } else {
1012 info!(
1013 run_id = %run_id,
1014 step_id = %step.id,
1015 step_name = %step.name,
1016 from = %step.status.state,
1017 to = %target_status,
1018 "cleaned up orphaned step"
1019 );
1020 }
1021 }
1022
1023 Ok(())
1024 }
1025
1026 async fn finalize_run(
1032 &self,
1033 run_id: Uuid,
1034 workflow_name: &str,
1035 result: Result<(), EngineError>,
1036 ctx: &WorkflowContext,
1037 run_start: Instant,
1038 run_labels: HashMap<String, String>,
1039 ) -> Result<Run, EngineError> {
1040 let total_duration = ctx.carried_duration_ms() + run_start.elapsed().as_millis() as u64;
1043 let completed_at = Utc::now();
1044
1045 let final_status;
1046 let final_run;
1047
1048 match result {
1049 Ok(()) => {
1050 final_status = if ctx.has_allowed_failure() {
1051 RunStatus::Warning
1052 } else {
1053 RunStatus::Completed
1054 };
1055 final_run = self
1056 .store
1057 .update_run_returning(
1058 run_id,
1059 RunUpdate {
1060 status: Some(final_status),
1061 cost_usd: Some(ctx.total_cost_usd()),
1062 duration_ms: Some(total_duration),
1063 completed_at: Some(completed_at),
1064 ..RunUpdate::default()
1065 },
1066 )
1067 .await?;
1068
1069 info!(
1070 run_id = %run_id,
1071 status = %final_status,
1072 cost_usd = %ctx.total_cost_usd(),
1073 duration_ms = total_duration,
1074 "run completed"
1075 );
1076 }
1077 Err(EngineError::ApprovalRequired {
1078 run_id: approval_run_id,
1079 step_id,
1080 ref message,
1081 }) => {
1082 final_status = RunStatus::AwaitingApproval;
1083 final_run = self
1084 .store
1085 .update_run_returning(
1086 run_id,
1087 RunUpdate {
1088 status: Some(RunStatus::AwaitingApproval),
1089 cost_usd: Some(ctx.total_cost_usd()),
1090 duration_ms: Some(total_duration),
1091 ..RunUpdate::default()
1092 },
1093 )
1094 .await?;
1095
1096 info!(
1097 run_id = %approval_run_id,
1098 step_id = %step_id,
1099 message = %message,
1100 "run awaiting approval"
1101 );
1102 }
1103 Err(err) => {
1104 let budget_exceeded = matches!(err, EngineError::RunBudgetExceeded { .. });
1108
1109 final_status = if budget_exceeded {
1110 if let Err(store_err) = self
1111 .store
1112 .update_run(
1113 run_id,
1114 RunUpdate {
1115 status: Some(RunStatus::Cancelled),
1116 error: Some(err.to_string()),
1117 cost_usd: Some(ctx.total_cost_usd()),
1118 duration_ms: Some(total_duration),
1119 completed_at: Some(completed_at),
1120 ..RunUpdate::default()
1121 },
1122 )
1123 .await
1124 {
1125 error!(run_id = %run_id, store_error = %store_err, "failed to persist run cancellation");
1126 }
1127 if let Err(cleanup_err) = self
1128 .fail_orphaned_steps(run_id, "run stopped: cost cap reached")
1129 .await
1130 {
1131 error!(run_id = %run_id, store_error = %cleanup_err, "failed to cleanup orphaned steps");
1132 }
1133 RunStatus::Cancelled
1134 } else {
1135 self.fail_or_schedule_retry(
1136 run_id,
1137 &err.to_string(),
1138 is_run_retryable(&err),
1139 Some(ctx.total_cost_usd()),
1140 Some(total_duration),
1141 )
1142 .await
1143 .unwrap_or_else(|store_err| {
1144 error!(run_id = %run_id, store_error = %store_err, "failed to persist run failure");
1145 RunStatus::Failed
1146 })
1147 };
1148
1149 if budget_exceeded {
1150 self.on_run_budget_exceeded(workflow_name, run_id, &err);
1151 }
1152
1153 error!(run_id = %run_id, status = %final_status, error = %err, "run stopped");
1154
1155 self.publish_run_status_changed(
1156 workflow_name,
1157 run_id,
1158 final_status,
1159 Some(err.to_string()),
1160 ctx,
1161 total_duration,
1162 run_labels,
1163 );
1164
1165 #[cfg(feature = "prometheus")]
1166 self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
1167
1168 return Err(err);
1169 }
1170 }
1171
1172 self.publish_run_status_changed(
1173 workflow_name,
1174 run_id,
1175 final_status,
1176 None,
1177 ctx,
1178 total_duration,
1179 run_labels,
1180 );
1181
1182 #[cfg(feature = "prometheus")]
1183 self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
1184
1185 Ok(final_run)
1186 }
1187
1188 #[cfg(feature = "prometheus")]
1190 fn emit_run_metrics(
1191 &self,
1192 workflow_name: &str,
1193 status: RunStatus,
1194 duration_ms: u64,
1195 ctx: &WorkflowContext,
1196 ) {
1197 let status_str = status.to_string();
1198 let wf = workflow_name.to_string();
1199
1200 counter!(RUNS_TOTAL, "workflow" => wf.clone(), "status" => status_str.clone()).increment(1);
1201 histogram!(RUN_DURATION_SECONDS, "workflow" => wf.clone(), "status" => status_str)
1202 .record(duration_ms as f64 / 1000.0);
1203 histogram!(RUN_COST_USD, "workflow" => wf.clone()).record(
1204 ctx.total_cost_usd()
1205 .to_string()
1206 .parse::<f64>()
1207 .unwrap_or(0.0),
1208 );
1209 gauge!(RUNS_ACTIVE, "workflow" => wf).decrement(1.0);
1210 }
1211
1212 fn on_run_budget_exceeded(&self, workflow_name: &str, run_id: Uuid, err: &EngineError) {
1218 let EngineError::RunBudgetExceeded {
1219 limit_usd,
1220 spent_usd,
1221 step_budget_usd,
1222 ..
1223 } = err
1224 else {
1225 return;
1226 };
1227
1228 #[cfg(feature = "prometheus")]
1229 counter!(
1230 RUN_BUDGET_EXCEEDED_TOTAL,
1231 "workflow" => workflow_name.to_string(),
1232 "scope" => "run",
1233 )
1234 .increment(1);
1235
1236 self.event_publisher.publish(Event::RunBudgetExceeded {
1237 run_id,
1238 workflow_name: workflow_name.to_string(),
1239 limit_usd: *limit_usd,
1240 spent_usd: *spent_usd,
1241 step_budget_usd: *step_budget_usd,
1242 at: Utc::now(),
1243 });
1244 }
1245
1246 #[allow(clippy::too_many_arguments)]
1251 fn publish_run_status_changed(
1252 &self,
1253 workflow_name: &str,
1254 run_id: Uuid,
1255 to: RunStatus,
1256 error: Option<String>,
1257 ctx: &WorkflowContext,
1258 duration_ms: u64,
1259 labels: HashMap<String, String>,
1260 ) {
1261 let now = Utc::now();
1262 let cost_usd = ctx.total_cost_usd();
1263 let wf = workflow_name.to_string();
1264
1265 self.event_publisher.publish(Event::RunStatusChanged {
1266 run_id,
1267 workflow_name: wf.clone(),
1268 from: RunStatus::Running,
1269 to,
1270 error: error.clone(),
1271 cost_usd,
1272 duration_ms,
1273 labels: labels.clone(),
1274 at: now,
1275 });
1276
1277 if to == RunStatus::Failed {
1278 self.event_publisher.publish(Event::RunFailed {
1279 run_id,
1280 workflow_name: wf,
1281 error,
1282 cost_usd,
1283 duration_ms,
1284 labels,
1285 at: now,
1286 });
1287 }
1288 }
1289}
1290
1291impl fmt::Debug for Engine {
1292 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1293 f.debug_struct("Engine")
1294 .field("handlers", &self.handlers.keys().collect::<Vec<_>>())
1295 .finish_non_exhaustive()
1296 }
1297}
1298
1299#[cfg(test)]
1300mod tests {
1301 use super::*;
1302 use crate::config::ShellConfig;
1303 use crate::handler::{HandlerFuture, WorkflowHandler};
1304 use ironflow_core::providers::claude::ClaudeCodeProvider;
1305 use ironflow_core::providers::record_replay::RecordReplayProvider;
1306 use ironflow_store::memory::InMemoryStore;
1307 use ironflow_store::models::StepStatus;
1308 use serde_json::json;
1309
1310 struct EchoWorkflow;
1312
1313 impl WorkflowHandler for EchoWorkflow {
1314 fn name(&self) -> &str {
1315 "echo-workflow"
1316 }
1317
1318 fn describe(&self) -> WorkflowInfo {
1319 WorkflowInfo {
1320 description: "A simple workflow that echoes hello".to_string(),
1321 source_code: None,
1322 sub_workflows: Vec::new(),
1323 category: None,
1324 version: self.version().map(str::to_string),
1325 compatible_versions: Vec::new(),
1326 input_schema: None,
1327 default_labels: HashMap::new(),
1328 schedule: self.schedule().cloned(),
1329 default_max_cost_usd: self.default_max_cost_usd(),
1330 }
1331 }
1332
1333 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
1334 Box::pin(async move {
1335 ctx.shell("greet", ShellConfig::new("echo hello")).await?;
1336 Ok(())
1337 })
1338 }
1339 }
1340
1341 struct FailingWorkflow;
1343
1344 impl WorkflowHandler for FailingWorkflow {
1345 fn name(&self) -> &str {
1346 "failing-workflow"
1347 }
1348
1349 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
1350 Box::pin(async move {
1351 ctx.shell("fail", ShellConfig::new("exit 1")).await?;
1352 Ok(())
1353 })
1354 }
1355 }
1356
1357 fn create_test_engine() -> Engine {
1358 let store = Arc::new(InMemoryStore::new());
1359 let inner = ClaudeCodeProvider::new();
1360 let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
1361 inner,
1362 "/tmp/ironflow-fixtures",
1363 ));
1364 Engine::new(store, provider)
1365 }
1366
1367 #[test]
1368 fn engine_new_creates_instance() {
1369 let engine = create_test_engine();
1370 assert_eq!(engine.handler_names().len(), 0);
1371 }
1372
1373 #[test]
1374 fn engine_register_handler() {
1375 let mut engine = create_test_engine();
1376 let result = engine.register(EchoWorkflow);
1377 assert!(result.is_ok());
1378 assert_eq!(engine.handler_names().len(), 1);
1379 assert!(engine.handler_names().contains(&"echo-workflow"));
1380 }
1381
1382 #[test]
1383 fn engine_register_duplicate_returns_error() {
1384 let mut engine = create_test_engine();
1385 engine.register(EchoWorkflow).unwrap();
1386 let result = engine.register(EchoWorkflow);
1387 assert!(result.is_err());
1388 }
1389
1390 #[test]
1391 fn engine_get_handler_found() {
1392 let mut engine = create_test_engine();
1393 engine.register(EchoWorkflow).unwrap();
1394 let handler = engine.get_handler("echo-workflow");
1395 assert!(handler.is_some());
1396 }
1397
1398 #[test]
1399 fn engine_get_handler_not_found() {
1400 let engine = create_test_engine();
1401 let handler = engine.get_handler("nonexistent");
1402 assert!(handler.is_none());
1403 }
1404
1405 #[test]
1406 fn engine_handler_names_lists_all() {
1407 let mut engine = create_test_engine();
1408 engine.register(EchoWorkflow).unwrap();
1409 engine.register(FailingWorkflow).unwrap();
1410 let names = engine.handler_names();
1411 assert_eq!(names.len(), 2);
1412 assert!(names.contains(&"echo-workflow"));
1413 assert!(names.contains(&"failing-workflow"));
1414 }
1415
1416 #[test]
1417 fn engine_handler_info_returns_description() {
1418 let mut engine = create_test_engine();
1419 engine.register(EchoWorkflow).unwrap();
1420 let info = engine.handler_info("echo-workflow");
1421 assert!(info.is_some());
1422 let info = info.unwrap();
1423 assert_eq!(info.description, "A simple workflow that echoes hello");
1424 }
1425
1426 struct CategorizedWorkflow;
1427
1428 impl WorkflowHandler for CategorizedWorkflow {
1429 fn name(&self) -> &str {
1430 "categorized"
1431 }
1432 fn category(&self) -> Option<&str> {
1433 Some("data/etl")
1434 }
1435 fn execute<'a>(
1436 &'a self,
1437 _ctx: &'a mut WorkflowContext,
1438 ) -> crate::handler::HandlerFuture<'a> {
1439 Box::pin(async move { Ok(()) })
1440 }
1441 }
1442
1443 #[test]
1444 fn engine_default_describe_propagates_category() {
1445 let mut engine = create_test_engine();
1446 engine.register(CategorizedWorkflow).unwrap();
1447 let info = engine.handler_info("categorized").unwrap();
1448 assert_eq!(info.category.as_deref(), Some("data/etl"));
1449 }
1450
1451 #[test]
1452 fn engine_default_describe_without_category() {
1453 let mut engine = create_test_engine();
1454 engine.register(EchoWorkflow).unwrap();
1455 let info = engine.handler_info("echo-workflow").unwrap();
1456 assert!(info.category.is_none());
1457 }
1458
1459 struct ScheduledWorkflow {
1464 schedule: CronSchedule,
1465 }
1466
1467 impl ScheduledWorkflow {
1468 fn new() -> Self {
1469 Self {
1470 schedule: CronSchedule::new("0 0 * * * *").unwrap(),
1471 }
1472 }
1473 }
1474
1475 impl WorkflowHandler for ScheduledWorkflow {
1476 fn name(&self) -> &str {
1477 "scheduled"
1478 }
1479 fn schedule(&self) -> Option<&CronSchedule> {
1480 Some(&self.schedule)
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_default_describe_propagates_schedule() {
1492 let mut engine = create_test_engine();
1493 engine.register(ScheduledWorkflow::new()).unwrap();
1494 let info = engine.handler_info("scheduled").unwrap();
1495 assert_eq!(
1496 info.schedule.as_ref().map(|s| s.as_str()),
1497 Some("0 0 * * * *")
1498 );
1499 }
1500
1501 #[test]
1502 fn engine_default_describe_without_schedule() {
1503 let mut engine = create_test_engine();
1504 engine.register(EchoWorkflow).unwrap();
1505 let info = engine.handler_info("echo-workflow").unwrap();
1506 assert!(info.schedule.is_none());
1507 }
1508
1509 #[test]
1510 fn scheduled_handlers_returns_only_scheduled() {
1511 let mut engine = create_test_engine();
1512 engine.register(EchoWorkflow).unwrap();
1513 engine.register(ScheduledWorkflow::new()).unwrap();
1514 engine.register(FailingWorkflow).unwrap();
1515
1516 let scheduled = engine.scheduled_handlers();
1517 assert_eq!(scheduled.len(), 1);
1518 assert_eq!(scheduled[0].0, "scheduled");
1519 assert_eq!(scheduled[0].1.as_str(), "0 0 * * * *");
1520 }
1521
1522 #[test]
1523 fn scheduled_handlers_empty_when_none_scheduled() {
1524 let mut engine = create_test_engine();
1525 engine.register(EchoWorkflow).unwrap();
1526 engine.register(FailingWorkflow).unwrap();
1527
1528 let scheduled = engine.scheduled_handlers();
1529 assert!(scheduled.is_empty());
1530 }
1531
1532 struct BadCategoryWorkflow(&'static str);
1533
1534 impl WorkflowHandler for BadCategoryWorkflow {
1535 fn name(&self) -> &str {
1536 "bad-category"
1537 }
1538 fn category(&self) -> Option<&str> {
1539 Some(self.0)
1540 }
1541 fn execute<'a>(
1542 &'a self,
1543 _ctx: &'a mut WorkflowContext,
1544 ) -> crate::handler::HandlerFuture<'a> {
1545 Box::pin(async move { Ok(()) })
1546 }
1547 }
1548
1549 #[test]
1550 fn engine_register_rejects_empty_category() {
1551 let mut engine = create_test_engine();
1552 let err = engine.register(BadCategoryWorkflow("")).unwrap_err();
1553 match err {
1554 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty category")),
1555 other => panic!("expected InvalidWorkflow, got {other:?}"),
1556 }
1557 }
1558
1559 #[test]
1560 fn engine_register_rejects_leading_slash_category() {
1561 let mut engine = create_test_engine();
1562 let err = engine
1563 .register(BadCategoryWorkflow("/data/etl"))
1564 .unwrap_err();
1565 match err {
1566 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("leading '/'")),
1567 other => panic!("expected InvalidWorkflow, got {other:?}"),
1568 }
1569 }
1570
1571 #[test]
1572 fn engine_register_rejects_trailing_slash_category() {
1573 let mut engine = create_test_engine();
1574 let err = engine
1575 .register(BadCategoryWorkflow("data/etl/"))
1576 .unwrap_err();
1577 match err {
1578 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("trailing '/'")),
1579 other => panic!("expected InvalidWorkflow, got {other:?}"),
1580 }
1581 }
1582
1583 #[test]
1584 fn engine_register_rejects_double_slash_category() {
1585 let mut engine = create_test_engine();
1586 let err = engine
1587 .register(BadCategoryWorkflow("data//etl"))
1588 .unwrap_err();
1589 match err {
1590 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty segment")),
1591 other => panic!("expected InvalidWorkflow, got {other:?}"),
1592 }
1593 }
1594
1595 #[test]
1596 fn engine_register_rejects_whitespace_only_segment_category() {
1597 let mut engine = create_test_engine();
1598 let err = engine
1599 .register(BadCategoryWorkflow("data/ /etl"))
1600 .unwrap_err();
1601 match err {
1602 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("whitespace-only segment")),
1603 other => panic!("expected InvalidWorkflow, got {other:?}"),
1604 }
1605 }
1606
1607 #[test]
1608 fn engine_register_accepts_valid_nested_category() {
1609 let mut engine = create_test_engine();
1610 assert!(engine.register(CategorizedWorkflow).is_ok());
1611 }
1612
1613 #[tokio::test]
1614 async fn engine_unknown_workflow_returns_error() {
1615 let engine = create_test_engine();
1616 let result = engine
1617 .run_handler("unknown", TriggerKind::Manual, json!({}))
1618 .await;
1619 assert!(result.is_err());
1620 match result {
1621 Err(EngineError::InvalidWorkflow(msg)) => {
1622 assert!(msg.contains("no handler registered"));
1623 }
1624 _ => panic!("expected InvalidWorkflow error"),
1625 }
1626 }
1627
1628 #[tokio::test]
1629 async fn engine_enqueue_handler_creates_pending_run() {
1630 let mut engine = create_test_engine();
1631 engine.register(EchoWorkflow).unwrap();
1632
1633 let run = engine
1634 .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
1635 .await
1636 .unwrap();
1637 assert_eq!(run.status.state, RunStatus::Pending);
1638 assert_eq!(run.workflow_name, "echo-workflow");
1639 }
1640
1641 #[tokio::test]
1642 async fn enqueue_handler_leaves_the_run_unattributed() {
1643 let mut engine = create_test_engine();
1644 engine.register(EchoWorkflow).unwrap();
1645
1646 let run = engine
1647 .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
1648 .await
1649 .unwrap();
1650
1651 assert!(run.created_by.is_none());
1652 }
1653
1654 #[tokio::test]
1655 async fn enqueue_handler_with_options_records_the_author() {
1656 let mut engine = create_test_engine();
1657 engine.register(EchoWorkflow).unwrap();
1658 let actor = RunActor::User {
1659 user_id: Uuid::now_v7(),
1660 };
1661
1662 let run = engine
1663 .enqueue_handler_with_options(
1664 "echo-workflow",
1665 TriggerKind::Api,
1666 json!({}),
1667 EnqueueOptions {
1668 created_by: Some(actor.clone()),
1669 ..Default::default()
1670 },
1671 )
1672 .await
1673 .unwrap()
1674 .into_run();
1675
1676 assert_eq!(run.created_by, Some(actor));
1677 }
1678
1679 #[tokio::test]
1680 async fn enqueue_handler_with_options_accepts_no_author() {
1681 let mut engine = create_test_engine();
1682 engine.register(EchoWorkflow).unwrap();
1683
1684 let run = engine
1685 .enqueue_handler_with_options(
1686 "echo-workflow",
1687 TriggerKind::Cron {
1688 schedule: "0 * * * * *".to_string(),
1689 },
1690 json!({}),
1691 EnqueueOptions::default(),
1692 )
1693 .await
1694 .unwrap()
1695 .into_run();
1696
1697 assert!(run.created_by.is_none());
1698 }
1699
1700 #[tokio::test]
1701 async fn run_handler_leaves_the_run_unattributed() {
1702 let mut engine = create_test_engine();
1703 engine.register(EchoWorkflow).unwrap();
1704
1705 let run = engine
1706 .run_handler("echo-workflow", TriggerKind::Manual, json!({}))
1707 .await
1708 .unwrap();
1709
1710 assert!(run.created_by.is_none());
1711 }
1712
1713 #[tokio::test]
1714 async fn engine_register_boxed() {
1715 let mut engine = create_test_engine();
1716 let handler: Box<dyn WorkflowHandler> = Box::new(EchoWorkflow);
1717 let result = engine.register_boxed(handler);
1718 assert!(result.is_ok());
1719 assert_eq!(engine.handler_names().len(), 1);
1720 }
1721
1722 #[tokio::test]
1723 async fn engine_store_and_provider_accessors() {
1724 let store = Arc::new(InMemoryStore::new());
1725 let inner = ClaudeCodeProvider::new();
1726 let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
1727 inner,
1728 "/tmp/ironflow-fixtures",
1729 ));
1730 let engine = Engine::new(store.clone(), provider.clone());
1731
1732 let _ = engine.store();
1734 let _ = engine.provider();
1735 }
1736
1737 use crate::operation::Operation;
1742 use ironflow_store::models::StepKind;
1743 use std::future::Future;
1744 use std::pin::Pin;
1745
1746 struct FakeGitlabOp {
1747 project_id: u64,
1748 title: String,
1749 }
1750
1751 impl Operation for FakeGitlabOp {
1752 fn kind(&self) -> &str {
1753 "gitlab"
1754 }
1755
1756 fn execute(&self) -> Pin<Box<dyn Future<Output = Result<Value, EngineError>> + Send + '_>> {
1757 Box::pin(async move {
1758 Ok(json!({
1759 "issue_id": 42,
1760 "project_id": self.project_id,
1761 "title": self.title,
1762 }))
1763 })
1764 }
1765
1766 fn input(&self) -> Option<Value> {
1767 Some(json!({
1768 "project_id": self.project_id,
1769 "title": self.title,
1770 }))
1771 }
1772 }
1773
1774 struct FailingOp;
1775
1776 impl Operation for FailingOp {
1777 fn kind(&self) -> &str {
1778 "broken-service"
1779 }
1780
1781 fn execute(&self) -> Pin<Box<dyn Future<Output = Result<Value, EngineError>> + Send + '_>> {
1782 Box::pin(async move { Err(EngineError::StepConfig("service unavailable".to_string())) })
1783 }
1784 }
1785
1786 struct OperationWorkflow;
1787
1788 impl WorkflowHandler for OperationWorkflow {
1789 fn name(&self) -> &str {
1790 "operation-workflow"
1791 }
1792
1793 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
1794 Box::pin(async move {
1795 let op = FakeGitlabOp {
1796 project_id: 123,
1797 title: "Bug report".to_string(),
1798 };
1799 ctx.operation("create-issue", &op).await?;
1800 Ok(())
1801 })
1802 }
1803 }
1804
1805 struct FailingOperationWorkflow;
1806
1807 impl WorkflowHandler for FailingOperationWorkflow {
1808 fn name(&self) -> &str {
1809 "failing-operation-workflow"
1810 }
1811
1812 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
1813 Box::pin(async move {
1814 ctx.operation("broken-call", &FailingOp).await?;
1815 Ok(())
1816 })
1817 }
1818 }
1819
1820 struct MixedWorkflow;
1821
1822 impl WorkflowHandler for MixedWorkflow {
1823 fn name(&self) -> &str {
1824 "mixed-workflow"
1825 }
1826
1827 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
1828 Box::pin(async move {
1829 ctx.shell("build", ShellConfig::new("echo built")).await?;
1830 let op = FakeGitlabOp {
1831 project_id: 456,
1832 title: "Deploy done".to_string(),
1833 };
1834 let result = ctx.operation("notify-gitlab", &op).await?;
1835 assert_eq!(result.output["issue_id"], 42);
1836 Ok(())
1837 })
1838 }
1839 }
1840
1841 #[tokio::test]
1842 async fn operation_step_happy_path() {
1843 let mut engine = create_test_engine();
1844 engine.register(OperationWorkflow).unwrap();
1845
1846 let run = engine
1847 .run_handler("operation-workflow", TriggerKind::Manual, json!({}))
1848 .await
1849 .unwrap();
1850
1851 assert_eq!(run.status.state, RunStatus::Completed);
1852
1853 let steps = engine.store().list_steps(run.id).await.unwrap();
1854
1855 assert_eq!(steps.len(), 1);
1856 assert_eq!(steps[0].name, "create-issue");
1857 assert_eq!(steps[0].kind, StepKind::Custom("gitlab".to_string()));
1858 assert_eq!(
1859 steps[0].status.state,
1860 ironflow_store::models::StepStatus::Completed
1861 );
1862
1863 let output = steps[0].output.as_ref().unwrap();
1864 assert_eq!(output["issue_id"], 42);
1865 assert_eq!(output["project_id"], 123);
1866
1867 let input = steps[0].input.as_ref().unwrap();
1868 assert_eq!(input["project_id"], 123);
1869 assert_eq!(input["title"], "Bug report");
1870 }
1871
1872 #[tokio::test]
1873 async fn operation_step_failure_marks_run_failed() {
1874 let mut engine = create_test_engine();
1875 engine.register(FailingOperationWorkflow).unwrap();
1876
1877 let result = engine
1878 .run_handler("failing-operation-workflow", TriggerKind::Manual, json!({}))
1879 .await;
1880
1881 assert!(result.is_err());
1882 }
1883
1884 #[tokio::test]
1885 async fn operation_mixed_with_shell_steps() {
1886 let mut engine = create_test_engine();
1887 engine.register(MixedWorkflow).unwrap();
1888
1889 let run = engine
1890 .run_handler("mixed-workflow", TriggerKind::Manual, json!({}))
1891 .await
1892 .unwrap();
1893
1894 assert_eq!(run.status.state, RunStatus::Completed);
1895
1896 let steps = engine.store().list_steps(run.id).await.unwrap();
1897
1898 assert_eq!(steps.len(), 2);
1899 assert_eq!(steps[0].kind, StepKind::Shell);
1900 assert_eq!(steps[1].kind, StepKind::Custom("gitlab".to_string()));
1901 assert_eq!(steps[0].position, 0);
1902 assert_eq!(steps[1].position, 1);
1903 }
1904
1905 use crate::config::ApprovalConfig;
1910
1911 struct SingleApprovalWorkflow;
1912
1913 impl WorkflowHandler for SingleApprovalWorkflow {
1914 fn name(&self) -> &str {
1915 "single-approval"
1916 }
1917
1918 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
1919 Box::pin(async move {
1920 ctx.shell("build", ShellConfig::new("echo built")).await?;
1921 ctx.approval("gate", ApprovalConfig::new("OK?")).await?;
1922 ctx.shell("deploy", ShellConfig::new("echo deployed"))
1923 .await?;
1924 Ok(())
1925 })
1926 }
1927 }
1928
1929 struct DoubleApprovalWorkflow;
1930
1931 impl WorkflowHandler for DoubleApprovalWorkflow {
1932 fn name(&self) -> &str {
1933 "double-approval"
1934 }
1935
1936 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
1937 Box::pin(async move {
1938 ctx.shell("build", ShellConfig::new("echo built")).await?;
1939 ctx.approval("staging-gate", ApprovalConfig::new("Deploy staging?"))
1940 .await?;
1941 ctx.shell("deploy-staging", ShellConfig::new("echo staging"))
1942 .await?;
1943 ctx.approval("prod-gate", ApprovalConfig::new("Deploy prod?"))
1944 .await?;
1945 ctx.shell("deploy-prod", ShellConfig::new("echo prod"))
1946 .await?;
1947 Ok(())
1948 })
1949 }
1950 }
1951
1952 #[tokio::test]
1953 async fn approval_pauses_run() {
1954 let mut engine = create_test_engine();
1955 engine.register(SingleApprovalWorkflow).unwrap();
1956
1957 let run = engine
1958 .run_handler("single-approval", TriggerKind::Manual, json!({}))
1959 .await
1960 .unwrap();
1961
1962 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
1963
1964 let steps = engine.store().list_steps(run.id).await.unwrap();
1965 assert_eq!(steps.len(), 2); assert_eq!(steps[0].kind, StepKind::Shell);
1967 assert_eq!(steps[0].status.state, StepStatus::Completed);
1968 assert_eq!(steps[1].kind, StepKind::Approval);
1969 assert_eq!(steps[1].status.state, StepStatus::AwaitingApproval);
1970 }
1971
1972 #[tokio::test]
1973 async fn approval_resume_completes_run() {
1974 let mut engine = create_test_engine();
1975 engine.register(SingleApprovalWorkflow).unwrap();
1976
1977 let run = engine
1979 .run_handler("single-approval", TriggerKind::Manual, json!({}))
1980 .await
1981 .unwrap();
1982 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
1983
1984 engine
1986 .store()
1987 .update_run_status(run.id, RunStatus::Running)
1988 .await
1989 .unwrap();
1990
1991 let resumed = engine.resume_run(run.id).await.unwrap();
1993 assert_eq!(resumed.status.state, RunStatus::Completed);
1994
1995 let steps = engine.store().list_steps(run.id).await.unwrap();
1996 assert_eq!(steps.len(), 3); assert_eq!(steps[0].name, "build");
1998 assert_eq!(steps[0].status.state, StepStatus::Completed);
1999 assert_eq!(steps[1].name, "gate");
2000 assert_eq!(steps[1].kind, StepKind::Approval);
2001 assert_eq!(steps[1].status.state, StepStatus::Completed);
2002 assert_eq!(steps[2].name, "deploy");
2003 assert_eq!(steps[2].status.state, StepStatus::Completed);
2004 }
2005
2006 #[tokio::test]
2007 async fn double_approval_two_resumes() {
2008 let mut engine = create_test_engine();
2009 engine.register(DoubleApprovalWorkflow).unwrap();
2010
2011 let run = engine
2013 .run_handler("double-approval", TriggerKind::Manual, json!({}))
2014 .await
2015 .unwrap();
2016 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
2017
2018 let steps = engine.store().list_steps(run.id).await.unwrap();
2019 assert_eq!(steps.len(), 2); engine
2023 .store()
2024 .update_run_status(run.id, RunStatus::Running)
2025 .await
2026 .unwrap();
2027
2028 let resumed = engine.resume_run(run.id).await.unwrap();
2029 assert_eq!(resumed.status.state, RunStatus::AwaitingApproval);
2030
2031 let steps = engine.store().list_steps(run.id).await.unwrap();
2032 assert_eq!(steps.len(), 4); engine
2036 .store()
2037 .update_run_status(run.id, RunStatus::Running)
2038 .await
2039 .unwrap();
2040
2041 let final_run = engine.resume_run(run.id).await.unwrap();
2042 assert_eq!(final_run.status.state, RunStatus::Completed);
2043
2044 let steps = engine.store().list_steps(run.id).await.unwrap();
2045 assert_eq!(steps.len(), 5);
2046 assert_eq!(steps[0].name, "build");
2047 assert_eq!(steps[1].name, "staging-gate");
2048 assert_eq!(steps[2].name, "deploy-staging");
2049 assert_eq!(steps[3].name, "prod-gate");
2050 assert_eq!(steps[4].name, "deploy-prod");
2051
2052 for step in &steps {
2053 assert_eq!(step.status.state, StepStatus::Completed);
2054 }
2055 }
2056
2057 use ironflow_store::models::{NewStep, StepUpdate};
2062
2063 async fn create_step_with_status(
2064 store: &Arc<dyn Store>,
2065 run_id: Uuid,
2066 name: &str,
2067 position: u32,
2068 status: StepStatus,
2069 ) -> ironflow_store::models::Step {
2070 let step = store
2071 .create_step(NewStep {
2072 run_id,
2073 name: name.to_string(),
2074 kind: StepKind::Shell,
2075 position,
2076 input: None,
2077 is_error_handler: false,
2078 })
2079 .await
2080 .unwrap();
2081
2082 match status {
2083 StepStatus::Pending => {}
2084 StepStatus::Running => {
2085 store
2086 .update_step(
2087 step.id,
2088 StepUpdate {
2089 status: Some(StepStatus::Running),
2090 ..StepUpdate::default()
2091 },
2092 )
2093 .await
2094 .unwrap();
2095 }
2096 StepStatus::Completed => {
2097 store
2098 .update_step(
2099 step.id,
2100 StepUpdate {
2101 status: Some(StepStatus::Running),
2102 ..StepUpdate::default()
2103 },
2104 )
2105 .await
2106 .unwrap();
2107 store
2108 .update_step(
2109 step.id,
2110 StepUpdate {
2111 status: Some(StepStatus::Completed),
2112 ..StepUpdate::default()
2113 },
2114 )
2115 .await
2116 .unwrap();
2117 }
2118 StepStatus::AwaitingApproval => {
2119 store
2120 .update_step(
2121 step.id,
2122 StepUpdate {
2123 status: Some(StepStatus::Running),
2124 ..StepUpdate::default()
2125 },
2126 )
2127 .await
2128 .unwrap();
2129 store
2130 .update_step(
2131 step.id,
2132 StepUpdate {
2133 status: Some(StepStatus::AwaitingApproval),
2134 ..StepUpdate::default()
2135 },
2136 )
2137 .await
2138 .unwrap();
2139 }
2140 _ => panic!("unsupported status for test helper: {status}"),
2141 }
2142
2143 store.get_step(step.id).await.unwrap().unwrap()
2144 }
2145
2146 #[tokio::test]
2147 async fn fail_orphaned_steps_marks_running_as_failed() {
2148 let engine = create_test_engine();
2149 let run = engine
2150 .store()
2151 .create_run(NewRun {
2152 created_by: None,
2153 workflow_name: "test".to_string(),
2154 trigger: TriggerKind::Manual,
2155 payload: json!({}),
2156 max_retries: 0,
2157 handler_version: None,
2158 labels: HashMap::new(),
2159 scheduled_at: None,
2160 idempotency_key: None,
2161 max_cost_usd: None,
2162 })
2163 .await
2164 .unwrap()
2165 .into_run();
2166
2167 let step = create_step_with_status(
2168 engine.store(),
2169 run.id,
2170 "running-step",
2171 0,
2172 StepStatus::Running,
2173 )
2174 .await;
2175
2176 engine
2177 .fail_orphaned_steps(run.id, "parent run timed out")
2178 .await
2179 .unwrap();
2180
2181 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
2182 assert_eq!(updated.status.state, StepStatus::Failed);
2183 assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
2184 assert!(updated.completed_at.is_some());
2185 }
2186
2187 #[tokio::test]
2188 async fn fail_orphaned_steps_marks_pending_as_skipped() {
2189 let engine = create_test_engine();
2190 let run = engine
2191 .store()
2192 .create_run(NewRun {
2193 created_by: None,
2194 workflow_name: "test".to_string(),
2195 trigger: TriggerKind::Manual,
2196 payload: json!({}),
2197 max_retries: 0,
2198 handler_version: None,
2199 labels: HashMap::new(),
2200 scheduled_at: None,
2201 idempotency_key: None,
2202 max_cost_usd: None,
2203 })
2204 .await
2205 .unwrap()
2206 .into_run();
2207
2208 let step = create_step_with_status(
2209 engine.store(),
2210 run.id,
2211 "pending-step",
2212 0,
2213 StepStatus::Pending,
2214 )
2215 .await;
2216
2217 engine
2218 .fail_orphaned_steps(run.id, "parent run timed out")
2219 .await
2220 .unwrap();
2221
2222 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
2223 assert_eq!(updated.status.state, StepStatus::Skipped);
2224 assert!(updated.error.is_none());
2225 assert!(updated.completed_at.is_some());
2226 }
2227
2228 #[tokio::test]
2229 async fn fail_orphaned_steps_marks_awaiting_approval_as_failed() {
2230 let engine = create_test_engine();
2231 let run = engine
2232 .store()
2233 .create_run(NewRun {
2234 created_by: None,
2235 workflow_name: "test".to_string(),
2236 trigger: TriggerKind::Manual,
2237 payload: json!({}),
2238 max_retries: 0,
2239 handler_version: None,
2240 labels: HashMap::new(),
2241 scheduled_at: None,
2242 idempotency_key: None,
2243 max_cost_usd: None,
2244 })
2245 .await
2246 .unwrap()
2247 .into_run();
2248
2249 let step = create_step_with_status(
2250 engine.store(),
2251 run.id,
2252 "approval-step",
2253 0,
2254 StepStatus::AwaitingApproval,
2255 )
2256 .await;
2257
2258 engine
2259 .fail_orphaned_steps(run.id, "parent run timed out")
2260 .await
2261 .unwrap();
2262
2263 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
2264 assert_eq!(updated.status.state, StepStatus::Failed);
2265 assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
2266 assert!(updated.completed_at.is_some());
2267 }
2268
2269 #[tokio::test]
2270 async fn fail_orphaned_steps_skips_terminal_steps() {
2271 let engine = create_test_engine();
2272 let run = engine
2273 .store()
2274 .create_run(NewRun {
2275 created_by: None,
2276 workflow_name: "test".to_string(),
2277 trigger: TriggerKind::Manual,
2278 payload: json!({}),
2279 max_retries: 0,
2280 handler_version: None,
2281 labels: HashMap::new(),
2282 scheduled_at: None,
2283 idempotency_key: None,
2284 max_cost_usd: None,
2285 })
2286 .await
2287 .unwrap()
2288 .into_run();
2289
2290 let completed_step =
2291 create_step_with_status(engine.store(), run.id, "done", 0, StepStatus::Completed).await;
2292 let running_step =
2293 create_step_with_status(engine.store(), run.id, "in-flight", 1, StepStatus::Running)
2294 .await;
2295
2296 engine
2297 .fail_orphaned_steps(run.id, "parent run timed out")
2298 .await
2299 .unwrap();
2300
2301 let completed = engine
2302 .store()
2303 .get_step(completed_step.id)
2304 .await
2305 .unwrap()
2306 .unwrap();
2307 assert_eq!(completed.status.state, StepStatus::Completed);
2308
2309 let failed = engine
2310 .store()
2311 .get_step(running_step.id)
2312 .await
2313 .unwrap()
2314 .unwrap();
2315 assert_eq!(failed.status.state, StepStatus::Failed);
2316 }
2317
2318 #[tokio::test]
2319 async fn fail_orphaned_steps_mixed_states() {
2320 let engine = create_test_engine();
2321 let run = engine
2322 .store()
2323 .create_run(NewRun {
2324 created_by: None,
2325 workflow_name: "test".to_string(),
2326 trigger: TriggerKind::Manual,
2327 payload: json!({}),
2328 max_retries: 0,
2329 handler_version: None,
2330 labels: HashMap::new(),
2331 scheduled_at: None,
2332 idempotency_key: None,
2333 max_cost_usd: None,
2334 })
2335 .await
2336 .unwrap()
2337 .into_run();
2338
2339 let s_completed =
2340 create_step_with_status(engine.store(), run.id, "step-1", 0, StepStatus::Completed)
2341 .await;
2342 let s_running =
2343 create_step_with_status(engine.store(), run.id, "step-2", 1, StepStatus::Running).await;
2344 let s_pending =
2345 create_step_with_status(engine.store(), run.id, "step-3", 2, StepStatus::Pending).await;
2346
2347 engine.fail_orphaned_steps(run.id, "timeout").await.unwrap();
2348
2349 let r_completed = engine
2350 .store()
2351 .get_step(s_completed.id)
2352 .await
2353 .unwrap()
2354 .unwrap();
2355 assert_eq!(r_completed.status.state, StepStatus::Completed);
2356
2357 let r_running = engine
2358 .store()
2359 .get_step(s_running.id)
2360 .await
2361 .unwrap()
2362 .unwrap();
2363 assert_eq!(r_running.status.state, StepStatus::Failed);
2364 assert_eq!(r_running.error.as_deref(), Some("timeout"));
2365
2366 let r_pending = engine
2367 .store()
2368 .get_step(s_pending.id)
2369 .await
2370 .unwrap()
2371 .unwrap();
2372 assert_eq!(r_pending.status.state, StepStatus::Skipped);
2373 assert!(r_pending.error.is_none());
2374 }
2375
2376 #[tokio::test]
2377 async fn fail_orphaned_steps_no_steps_is_noop() {
2378 let engine = create_test_engine();
2379 let run = engine
2380 .store()
2381 .create_run(NewRun {
2382 created_by: None,
2383 workflow_name: "test".to_string(),
2384 trigger: TriggerKind::Manual,
2385 payload: json!({}),
2386 max_retries: 0,
2387 handler_version: None,
2388 labels: HashMap::new(),
2389 scheduled_at: None,
2390 idempotency_key: None,
2391 max_cost_usd: None,
2392 })
2393 .await
2394 .unwrap()
2395 .into_run();
2396
2397 let result = engine.fail_orphaned_steps(run.id, "timeout").await;
2398 assert!(result.is_ok());
2399 }
2400
2401 #[tokio::test]
2402 async fn fail_orphaned_steps_preserves_existing_error() {
2403 let engine = create_test_engine();
2404 let run = engine
2405 .store()
2406 .create_run(NewRun {
2407 created_by: None,
2408 workflow_name: "test".to_string(),
2409 trigger: TriggerKind::Manual,
2410 payload: json!({}),
2411 max_retries: 0,
2412 handler_version: None,
2413 labels: HashMap::new(),
2414 scheduled_at: None,
2415 idempotency_key: None,
2416 max_cost_usd: None,
2417 })
2418 .await
2419 .unwrap()
2420 .into_run();
2421
2422 let step_with_error = create_step_with_status(
2423 engine.store(),
2424 run.id,
2425 "already-errored",
2426 0,
2427 StepStatus::Running,
2428 )
2429 .await;
2430
2431 engine
2432 .store()
2433 .update_step(
2434 step_with_error.id,
2435 StepUpdate {
2436 error: Some("real error from provider".to_string()),
2437 ..StepUpdate::default()
2438 },
2439 )
2440 .await
2441 .unwrap();
2442
2443 let step_no_error = create_step_with_status(
2444 engine.store(),
2445 run.id,
2446 "no-error-yet",
2447 1,
2448 StepStatus::Running,
2449 )
2450 .await;
2451
2452 engine
2453 .fail_orphaned_steps(run.id, "parent run failed")
2454 .await
2455 .unwrap();
2456
2457 let updated_with = engine
2458 .store()
2459 .get_step(step_with_error.id)
2460 .await
2461 .unwrap()
2462 .unwrap();
2463 assert_eq!(updated_with.status.state, StepStatus::Failed);
2464 assert_eq!(
2465 updated_with.error.as_deref(),
2466 Some("real error from provider"),
2467 );
2468
2469 let updated_without = engine
2470 .store()
2471 .get_step(step_no_error.id)
2472 .await
2473 .unwrap()
2474 .unwrap();
2475 assert_eq!(updated_without.status.state, StepStatus::Failed);
2476 assert_eq!(updated_without.error.as_deref(), Some("parent run failed"),);
2477 }
2478}