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::executor::StepResult;
39use crate::guard::{WorkflowGuardConfig, new_shared_guard_state};
40use crate::handler::{WorkflowHandler, WorkflowInfo};
41use crate::log_sender::LogSender;
42use crate::notify::{
43 Event, EventPublisher, EventSubscriber, RunBudgetExceededEvent, RunFailedEvent,
44 RunStatusChangedEvent, WorkflowEventBus,
45};
46use crate::retry_policy::{backoff_for_retry, is_run_retryable};
47use crate::schedule::CronSchedule;
48use ironflow_core::decision::DecisionProvider;
49
50#[derive(Debug, Clone)]
69pub struct WorkflowResult {
70 pub run: Run,
72 pub steps: Vec<StepResult>,
74}
75
76#[derive(Debug, Clone, Default)]
95pub struct EnqueueOptions {
96 pub max_retries: u32,
98 pub labels: HashMap<String, String>,
100 pub scheduled_at: Option<DateTime<Utc>>,
103 pub max_cost_usd: Option<Decimal>,
107 pub created_by: Option<RunActor>,
110 pub idempotency_key: Option<String>,
116}
117
118pub struct Engine {
159 store: Arc<dyn Store>,
160 provider: Arc<dyn AgentProvider>,
161 handlers: HashMap<String, Arc<dyn WorkflowHandler>>,
162 event_publisher: EventPublisher,
163 log_sender: Option<LogSender>,
164 budget: BudgetConfig,
165 artifact_sink: Option<Arc<dyn ArtifactSink>>,
166 guard_config: Option<WorkflowGuardConfig>,
167 event_bus: Option<WorkflowEventBus>,
168 decision_provider: Option<Arc<dyn DecisionProvider>>,
169}
170
171fn validate_category(handler_name: &str, category: &str) -> Result<(), EngineError> {
181 let reject = |reason: &str| {
182 Err(EngineError::InvalidWorkflow(format!(
183 "handler '{handler_name}' has invalid category '{category}': {reason}"
184 )))
185 };
186
187 if category.is_empty() {
188 return reject("empty category");
189 }
190 if category.starts_with('/') {
191 return reject("leading '/'");
192 }
193 if category.ends_with('/') {
194 return reject("trailing '/'");
195 }
196 for segment in category.split('/') {
197 if segment.is_empty() {
198 return reject("empty segment (double '/')");
199 }
200 if segment.trim().is_empty() {
201 return reject("whitespace-only segment");
202 }
203 }
204 Ok(())
205}
206
207impl Engine {
208 pub fn new(store: Arc<dyn Store>, provider: Arc<dyn AgentProvider>) -> Self {
224 Self {
225 store,
226 provider,
227 handlers: HashMap::new(),
228 event_publisher: EventPublisher::new(),
229 log_sender: None,
230 budget: BudgetConfig::new(),
231 artifact_sink: None,
232 guard_config: None,
233 event_bus: None,
234 decision_provider: None,
235 }
236 }
237
238 pub fn with_decision_provider(mut self, provider: Arc<dyn DecisionProvider>) -> Self {
259 self.decision_provider = Some(provider);
260 self
261 }
262
263 pub fn with_budget_config(mut self, budget: BudgetConfig) -> Self {
284 self.budget = budget;
285 self
286 }
287
288 pub fn budget_config(&self) -> &BudgetConfig {
290 &self.budget
291 }
292
293 pub fn with_guard_config(mut self, config: WorkflowGuardConfig) -> Self {
315 self.guard_config = Some(config);
316 self
317 }
318
319 pub fn guard_config(&self) -> Option<&WorkflowGuardConfig> {
321 self.guard_config.as_ref()
322 }
323
324 pub fn set_log_sender(&mut self, sender: LogSender) {
330 self.log_sender = Some(sender);
331 }
332
333 pub fn set_artifact_sink(&mut self, sink: Arc<dyn ArtifactSink>) {
351 self.artifact_sink = Some(sink);
352 }
353
354 pub fn artifact_sink(&self) -> Option<&Arc<dyn ArtifactSink>> {
356 self.artifact_sink.as_ref()
357 }
358
359 pub fn set_event_bus(&mut self, bus: WorkflowEventBus) {
376 self.event_bus = Some(bus);
377 }
378
379 pub fn event_bus(&self) -> Option<&WorkflowEventBus> {
381 self.event_bus.as_ref()
382 }
383
384 pub fn store(&self) -> &Arc<dyn Store> {
386 &self.store
387 }
388
389 pub fn provider(&self) -> &Arc<dyn AgentProvider> {
391 &self.provider
392 }
393
394 fn build_context(&self, run: &Run) -> WorkflowContext {
403 let handlers = self.handlers.clone();
404 let resolver: crate::context::HandlerResolver =
405 Arc::new(move |name: &str| handlers.get(name).cloned());
406 let mut ctx = WorkflowContext::with_handler_resolver(
407 run.id,
408 run.workflow_name.clone(),
409 self.store.clone(),
410 self.provider.clone(),
411 resolver,
412 );
413 ctx.carry_over_run_totals(run.retry_count + 1, run.cost_usd, run.duration_ms);
414 ctx.set_max_cost_usd(run.max_cost_usd);
415 if let Some(ref sender) = self.log_sender {
416 ctx.set_log_sender(sender.clone());
417 }
418 if let Some(ref sink) = self.artifact_sink {
419 ctx.set_artifact_sink(sink.clone());
420 }
421 if let Some(ref bus) = self.event_bus {
422 ctx.set_event_bus(bus.clone());
423 }
424 if let Some(ref provider) = self.decision_provider {
425 ctx.set_decision_provider(provider.clone());
426 }
427 ctx
428 }
429
430 fn build_context_with_guard(
436 &self,
437 run: &Run,
438 handler: &dyn WorkflowHandler,
439 ) -> WorkflowContext {
440 let mut ctx = self.build_context(run);
441 let guard_config = handler.guard_config().or_else(|| self.guard_config.clone());
442 if let Some(config) = guard_config {
443 ctx.set_guard(config, new_shared_guard_state());
444 }
445 ctx
446 }
447
448 async fn check_monthly_quota(&self, workflow_name: &str) -> Result<(), EngineError> {
459 let Some(limit) = self.budget.monthly_cost_limit_usd else {
460 return Ok(());
461 };
462
463 let stats = self
464 .store
465 .get_stats(RunFilter {
466 created_after: Some(month_start(Utc::now())),
467 ..RunFilter::default()
468 })
469 .await?;
470
471 if stats.total_cost_usd < limit {
472 return Ok(());
473 }
474
475 warn!(
476 workflow = %workflow_name,
477 limit_usd = %limit,
478 spent_usd = %stats.total_cost_usd,
479 "monthly cost quota exhausted, refusing new run"
480 );
481
482 #[cfg(feature = "prometheus")]
483 counter!(
484 RUN_BUDGET_EXCEEDED_TOTAL,
485 "workflow" => workflow_name.to_string(),
486 "scope" => "monthly",
487 )
488 .increment(1);
489
490 Err(EngineError::MonthlyBudgetExceeded {
491 limit_usd: limit,
492 spent_usd: stats.total_cost_usd,
493 })
494 }
495
496 pub fn register(&mut self, handler: impl WorkflowHandler + 'static) -> Result<(), EngineError> {
540 let name = handler.name().to_string();
541 if self.handlers.contains_key(&name) {
542 return Err(EngineError::InvalidWorkflow(format!(
543 "handler '{}' already registered",
544 name
545 )));
546 }
547 if let Some(category) = handler.category() {
548 validate_category(&name, category)?;
549 }
550 self.handlers.insert(name, Arc::new(handler));
551 Ok(())
552 }
553
554 pub fn register_boxed(&mut self, handler: Box<dyn WorkflowHandler>) -> Result<(), EngineError> {
561 let name = handler.name().to_string();
562 if self.handlers.contains_key(&name) {
563 return Err(EngineError::InvalidWorkflow(format!(
564 "handler '{}' already registered",
565 name
566 )));
567 }
568 if let Some(category) = handler.category() {
569 validate_category(&name, category)?;
570 }
571 self.handlers.insert(name, Arc::from(handler));
572 Ok(())
573 }
574
575 pub fn get_handler(&self, name: &str) -> Option<&Arc<dyn WorkflowHandler>> {
577 self.handlers.get(name)
578 }
579
580 pub fn handler_names(&self) -> Vec<&str> {
582 self.handlers.keys().map(|s| s.as_str()).collect()
583 }
584
585 pub fn handler_info(&self, name: &str) -> Option<WorkflowInfo> {
587 self.handlers.get(name).map(|h| h.describe())
588 }
589
590 pub fn scheduled_handlers(&self) -> Vec<(&str, &CronSchedule)> {
614 self.handlers
615 .iter()
616 .filter_map(|(name, handler)| handler.schedule().map(|sched| (name.as_str(), sched)))
617 .collect()
618 }
619
620 pub fn subscribe(
645 &mut self,
646 subscriber: impl EventSubscriber + 'static,
647 event_types: &[&'static str],
648 ) {
649 self.event_publisher.subscribe(subscriber, event_types);
650 }
651
652 pub fn event_publisher(&self) -> &EventPublisher {
657 &self.event_publisher
658 }
659
660 #[tracing::instrument(name = "engine.run_handler", skip_all, fields(workflow = %handler_name))]
690 pub async fn run_handler(
691 &self,
692 handler_name: &str,
693 trigger: TriggerKind,
694 payload: Value,
695 ) -> Result<WorkflowResult, EngineError> {
696 let handler = self
697 .handlers
698 .get(handler_name)
699 .ok_or_else(|| {
700 EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
701 })?
702 .clone();
703
704 self.check_monthly_quota(handler_name).await?;
705
706 let handler_version = handler.version().map(str::to_string);
707 let max_cost_usd = self
708 .budget
709 .resolve_run_cap(None, handler.default_max_cost_usd());
710 let run = self
711 .store
712 .create_run(NewRun {
713 created_by: None,
714 workflow_name: handler_name.to_string(),
715 trigger,
716 payload,
717 max_retries: 0,
718 handler_version,
719 labels: handler.default_labels(),
720 scheduled_at: None,
721 idempotency_key: None,
722 max_cost_usd,
723 })
724 .await?
725 .into_run();
726
727 let run_id = run.id;
728 info!(run_id = %run_id, handler_version = run.handler_version.as_deref().unwrap_or(""), "run created");
729
730 self.store
731 .update_run_status(run_id, RunStatus::Running)
732 .await?;
733
734 #[cfg(feature = "prometheus")]
735 gauge!(RUNS_ACTIVE, "workflow" => handler_name.to_string()).increment(1.0);
736
737 let run_start = Instant::now();
738 let mut ctx = self.build_context_with_guard(&run, handler.as_ref());
739
740 let result = handler.execute(&mut ctx).await;
741 self.finalize_run(run_id, handler_name, result, &ctx, run_start, run.labels)
742 .await
743 }
744
745 #[tracing::instrument(name = "engine.enqueue_handler", skip_all, fields(workflow = %handler_name))]
756 pub async fn enqueue_handler(
757 &self,
758 handler_name: &str,
759 trigger: TriggerKind,
760 payload: Value,
761 max_retries: u32,
762 ) -> Result<Run, EngineError> {
763 self.enqueue_handler_with_options(
764 handler_name,
765 trigger,
766 payload,
767 EnqueueOptions {
768 max_retries,
769 ..Default::default()
770 },
771 )
772 .await
773 .map(RunCreation::into_run)
774 }
775
776 #[tracing::instrument(name = "engine.enqueue_handler_with_options", skip_all, fields(workflow = %handler_name))]
821 pub async fn enqueue_handler_with_options(
822 &self,
823 handler_name: &str,
824 trigger: TriggerKind,
825 payload: Value,
826 options: EnqueueOptions,
827 ) -> Result<RunCreation, EngineError> {
828 let EnqueueOptions {
829 max_retries,
830 labels,
831 scheduled_at,
832 max_cost_usd,
833 created_by,
834 idempotency_key,
835 } = options;
836
837 let handler = self.handlers.get(handler_name).ok_or_else(|| {
838 EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
839 })?;
840
841 self.check_monthly_quota(handler_name).await?;
842
843 let handler_version = handler.version().map(str::to_string);
844 let mut merged_labels = handler.default_labels();
845 merged_labels.extend(labels);
846 let resolved_cap = self
847 .budget
848 .resolve_run_cap(max_cost_usd, handler.default_max_cost_usd());
849
850 let creation = self
851 .store
852 .create_run(NewRun {
853 workflow_name: handler_name.to_string(),
854 trigger,
855 payload,
856 max_retries,
857 handler_version,
858 labels: merged_labels,
859 scheduled_at,
860 created_by,
861 idempotency_key,
862 max_cost_usd: resolved_cap,
863 })
864 .await?;
865
866 match &creation {
867 RunCreation::Created(run) => info!(
868 run_id = %run.id,
869 workflow = %handler_name,
870 max_cost_usd = ?resolved_cap,
871 "handler run enqueued"
872 ),
873 RunCreation::Existing(run) => info!(
874 run_id = %run.id,
875 workflow = %handler_name,
876 "idempotent replay, nothing enqueued"
877 ),
878 }
879
880 Ok(creation)
881 }
882
883 #[tracing::instrument(name = "engine.execute_handler_run", skip_all, fields(run_id = %run_id))]
892 pub async fn execute_handler_run(&self, run_id: Uuid) -> Result<WorkflowResult, EngineError> {
893 let run = self
894 .store
895 .get_run(run_id)
896 .await?
897 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
898
899 let handler = self
900 .handlers
901 .get(&run.workflow_name)
902 .ok_or_else(|| {
903 EngineError::InvalidWorkflow(format!(
904 "no handler registered: {}",
905 run.workflow_name
906 ))
907 })?
908 .clone();
909
910 #[cfg(feature = "prometheus")]
911 gauge!(RUNS_ACTIVE, "workflow" => run.workflow_name.clone()).increment(1.0);
912
913 let run_start = Instant::now();
914 let mut ctx = self.build_context_with_guard(&run, handler.as_ref());
915
916 if run.retry_count > 0 {
919 ctx.load_replay_steps().await?;
920 }
921
922 let result = handler.execute(&mut ctx).await;
923 self.finalize_run(
924 run_id,
925 &run.workflow_name,
926 result,
927 &ctx,
928 run_start,
929 run.labels,
930 )
931 .await
932 }
933
934 #[tracing::instrument(name = "engine.execute_run", skip_all, fields(run_id = %run_id))]
942 pub async fn execute_run(&self, run_id: Uuid) -> Result<WorkflowResult, EngineError> {
943 self.execute_handler_run(run_id).await
944 }
945
946 #[tracing::instrument(name = "engine.resume_run", skip_all, fields(run_id = %run_id))]
960 pub async fn resume_run(&self, run_id: Uuid) -> Result<WorkflowResult, EngineError> {
961 let run = self
962 .store
963 .get_run(run_id)
964 .await?
965 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
966
967 let handler = self
968 .handlers
969 .get(&run.workflow_name)
970 .ok_or_else(|| {
971 EngineError::InvalidWorkflow(format!(
972 "no handler registered: {}",
973 run.workflow_name
974 ))
975 })?
976 .clone();
977
978 info!(run_id = %run_id, workflow = %run.workflow_name, "resuming run after approval");
979
980 let run_start = Instant::now();
981 let mut ctx = self.build_context_with_guard(&run, handler.as_ref());
982 ctx.load_replay_steps().await?;
983
984 let result = handler.execute(&mut ctx).await;
985 self.finalize_run(
986 run_id,
987 &run.workflow_name,
988 result,
989 &ctx,
990 run_start,
991 run.labels,
992 )
993 .await
994 }
995
996 pub async fn fail_or_schedule_retry(
1041 &self,
1042 run_id: Uuid,
1043 error: &str,
1044 retryable: bool,
1045 cost_usd: Option<Decimal>,
1046 duration_ms: Option<u64>,
1047 ) -> Result<RunStatus, EngineError> {
1048 let run = self
1049 .store
1050 .get_run(run_id)
1051 .await?
1052 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
1053
1054 let has_attempts_left = run.retry_count < run.max_retries;
1055 let update = if retryable && has_attempts_left {
1056 let backoff = backoff_for_retry(run.retry_count);
1057 let scheduled_at = Utc::now() + TimeDelta::milliseconds(backoff.as_millis() as i64);
1058
1059 info!(
1060 run_id = %run_id,
1061 workflow = %run.workflow_name,
1062 attempt = run.retry_count + 1,
1063 max_retries = run.max_retries,
1064 backoff_secs = backoff.as_secs(),
1065 scheduled_at = %scheduled_at,
1066 "run failed, scheduling retry"
1067 );
1068
1069 RunUpdate {
1070 status: Some(RunStatus::Retrying),
1071 error: Some(error.to_string()),
1072 increment_retry: true,
1073 cost_usd,
1074 duration_ms,
1075 scheduled_at: Some(scheduled_at),
1076 ..RunUpdate::default()
1077 }
1078 } else {
1079 RunUpdate {
1080 status: Some(RunStatus::Failed),
1081 error: Some(error.to_string()),
1082 cost_usd,
1083 duration_ms,
1084 completed_at: Some(Utc::now()),
1085 ..RunUpdate::default()
1086 }
1087 };
1088
1089 let status = update.status.unwrap_or(RunStatus::Failed);
1090 self.store.update_run(run_id, update).await?;
1091 self.fail_orphaned_steps(run_id, error).await?;
1092
1093 Ok(status)
1094 }
1095
1096 pub async fn fail_orphaned_steps(
1110 &self,
1111 run_id: Uuid,
1112 error_message: &str,
1113 ) -> Result<(), EngineError> {
1114 let steps = self.store.list_steps(run_id).await?;
1115 let now = Utc::now();
1116
1117 for step in steps {
1118 if step.status.state.is_terminal() {
1119 continue;
1120 }
1121
1122 let (target_status, error) = match step.status.state {
1123 StepStatus::Running | StepStatus::AwaitingApproval => {
1124 let err = if step.error.is_some() {
1125 None
1126 } else {
1127 Some(error_message.to_string())
1128 };
1129 (StepStatus::Failed, err)
1130 }
1131 StepStatus::Pending => (StepStatus::Skipped, None),
1132 _ => continue,
1133 };
1134
1135 if let Err(e) = self
1136 .store
1137 .update_step(
1138 step.id,
1139 StepUpdate {
1140 status: Some(target_status),
1141 error,
1142 completed_at: Some(now),
1143 ..StepUpdate::default()
1144 },
1145 )
1146 .await
1147 {
1148 warn!(
1149 run_id = %run_id,
1150 step_id = %step.id,
1151 step_name = %step.name,
1152 error = %e,
1153 "failed to cleanup orphaned step"
1154 );
1155 } else {
1156 info!(
1157 run_id = %run_id,
1158 step_id = %step.id,
1159 step_name = %step.name,
1160 from = %step.status.state,
1161 to = %target_status,
1162 "cleaned up orphaned step"
1163 );
1164 }
1165 }
1166
1167 Ok(())
1168 }
1169
1170 async fn finalize_run(
1176 &self,
1177 run_id: Uuid,
1178 workflow_name: &str,
1179 result: Result<(), EngineError>,
1180 ctx: &WorkflowContext,
1181 run_start: Instant,
1182 run_labels: HashMap<String, String>,
1183 ) -> Result<WorkflowResult, EngineError> {
1184 let total_duration = ctx.carried_duration_ms() + run_start.elapsed().as_millis() as u64;
1187 let completed_at = Utc::now();
1188
1189 let final_status;
1190 let final_run;
1191
1192 match result {
1193 Ok(()) => {
1194 final_status = if ctx.has_allowed_failure() {
1195 RunStatus::Warning
1196 } else {
1197 RunStatus::Completed
1198 };
1199 final_run = self
1200 .store
1201 .update_run_returning(
1202 run_id,
1203 RunUpdate {
1204 status: Some(final_status),
1205 cost_usd: Some(ctx.total_cost_usd()),
1206 duration_ms: Some(total_duration),
1207 completed_at: Some(completed_at),
1208 ..RunUpdate::default()
1209 },
1210 )
1211 .await?;
1212
1213 info!(
1214 run_id = %run_id,
1215 status = %final_status,
1216 cost_usd = %ctx.total_cost_usd(),
1217 duration_ms = total_duration,
1218 "run completed"
1219 );
1220 }
1221 Err(EngineError::ApprovalRequired {
1222 run_id: approval_run_id,
1223 step_id,
1224 ref message,
1225 }) => {
1226 final_status = RunStatus::AwaitingApproval;
1227 final_run = self
1228 .store
1229 .update_run_returning(
1230 run_id,
1231 RunUpdate {
1232 status: Some(RunStatus::AwaitingApproval),
1233 cost_usd: Some(ctx.total_cost_usd()),
1234 duration_ms: Some(total_duration),
1235 ..RunUpdate::default()
1236 },
1237 )
1238 .await?;
1239
1240 info!(
1241 run_id = %approval_run_id,
1242 step_id = %step_id,
1243 message = %message,
1244 "run awaiting approval"
1245 );
1246 }
1247 Err(EngineError::DelaySleeping {
1248 run_id: delay_run_id,
1249 step_id,
1250 wake_at,
1251 }) => {
1252 final_status = RunStatus::Sleeping;
1253 final_run = self
1254 .store
1255 .update_run_returning(
1256 run_id,
1257 RunUpdate {
1258 status: Some(RunStatus::Sleeping),
1259 cost_usd: Some(ctx.total_cost_usd()),
1260 duration_ms: Some(total_duration),
1261 scheduled_at: Some(wake_at),
1262 ..RunUpdate::default()
1263 },
1264 )
1265 .await?;
1266
1267 info!(
1268 run_id = %delay_run_id,
1269 step_id = %step_id,
1270 wake_at = %wake_at,
1271 "run sleeping until delay elapses"
1272 );
1273 }
1274 Err(err) => {
1275 let guardrail_stop = matches!(
1279 err,
1280 EngineError::RunBudgetExceeded { .. } | EngineError::WorkflowGuardRejected(_)
1281 );
1282
1283 final_status = if guardrail_stop {
1284 if let Err(store_err) = self
1285 .store
1286 .update_run(
1287 run_id,
1288 RunUpdate {
1289 status: Some(RunStatus::Cancelled),
1290 error: Some(err.to_string()),
1291 cost_usd: Some(ctx.total_cost_usd()),
1292 duration_ms: Some(total_duration),
1293 completed_at: Some(completed_at),
1294 ..RunUpdate::default()
1295 },
1296 )
1297 .await
1298 {
1299 error!(run_id = %run_id, store_error = %store_err, "failed to persist run cancellation");
1300 }
1301 if let Err(cleanup_err) = self
1302 .fail_orphaned_steps(run_id, "run stopped: guardrail limit reached")
1303 .await
1304 {
1305 error!(run_id = %run_id, store_error = %cleanup_err, "failed to cleanup orphaned steps");
1306 }
1307 RunStatus::Cancelled
1308 } else {
1309 self.fail_or_schedule_retry(
1310 run_id,
1311 &err.to_string(),
1312 is_run_retryable(&err),
1313 Some(ctx.total_cost_usd()),
1314 Some(total_duration),
1315 )
1316 .await
1317 .unwrap_or_else(|store_err| {
1318 error!(run_id = %run_id, store_error = %store_err, "failed to persist run failure");
1319 RunStatus::Failed
1320 })
1321 };
1322
1323 if matches!(err, EngineError::RunBudgetExceeded { .. }) {
1324 self.on_run_budget_exceeded(workflow_name, run_id, &err);
1325 }
1326
1327 error!(run_id = %run_id, status = %final_status, error = %err, "run stopped");
1328
1329 self.publish_run_status_changed(
1330 workflow_name,
1331 run_id,
1332 final_status,
1333 Some(err.to_string()),
1334 ctx,
1335 total_duration,
1336 run_labels,
1337 );
1338
1339 #[cfg(feature = "prometheus")]
1340 self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
1341
1342 return Err(err);
1343 }
1344 }
1345
1346 self.publish_run_status_changed(
1347 workflow_name,
1348 run_id,
1349 final_status,
1350 None,
1351 ctx,
1352 total_duration,
1353 run_labels,
1354 );
1355
1356 #[cfg(feature = "prometheus")]
1357 self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
1358
1359 Ok(WorkflowResult {
1360 run: final_run,
1361 steps: ctx.step_results().to_vec(),
1362 })
1363 }
1364
1365 #[cfg(feature = "prometheus")]
1367 fn emit_run_metrics(
1368 &self,
1369 workflow_name: &str,
1370 status: RunStatus,
1371 duration_ms: u64,
1372 ctx: &WorkflowContext,
1373 ) {
1374 let status_str = status.to_string();
1375 let wf = workflow_name.to_string();
1376
1377 counter!(RUNS_TOTAL, "workflow" => wf.clone(), "status" => status_str.clone()).increment(1);
1378 histogram!(RUN_DURATION_SECONDS, "workflow" => wf.clone(), "status" => status_str)
1379 .record(duration_ms as f64 / 1000.0);
1380 histogram!(RUN_COST_USD, "workflow" => wf.clone()).record(
1381 ctx.total_cost_usd()
1382 .to_string()
1383 .parse::<f64>()
1384 .unwrap_or(0.0),
1385 );
1386 gauge!(RUNS_ACTIVE, "workflow" => wf).decrement(1.0);
1387 }
1388
1389 fn on_run_budget_exceeded(&self, workflow_name: &str, run_id: Uuid, err: &EngineError) {
1395 let EngineError::RunBudgetExceeded {
1396 limit_usd,
1397 spent_usd,
1398 step_budget_usd,
1399 ..
1400 } = err
1401 else {
1402 return;
1403 };
1404
1405 #[cfg(feature = "prometheus")]
1406 counter!(
1407 RUN_BUDGET_EXCEEDED_TOTAL,
1408 "workflow" => workflow_name.to_string(),
1409 "scope" => "run",
1410 )
1411 .increment(1);
1412
1413 self.event_publisher
1414 .publish(Event::RunBudgetExceeded(RunBudgetExceededEvent {
1415 run_id,
1416 workflow_name: workflow_name.to_string(),
1417 limit_usd: *limit_usd,
1418 spent_usd: *spent_usd,
1419 step_budget_usd: *step_budget_usd,
1420 at: Utc::now(),
1421 }));
1422 }
1423
1424 #[allow(clippy::too_many_arguments)]
1429 fn publish_run_status_changed(
1430 &self,
1431 workflow_name: &str,
1432 run_id: Uuid,
1433 to: RunStatus,
1434 error: Option<String>,
1435 ctx: &WorkflowContext,
1436 duration_ms: u64,
1437 labels: HashMap<String, String>,
1438 ) {
1439 let now = Utc::now();
1440 let cost_usd = ctx.total_cost_usd();
1441 let wf = workflow_name.to_string();
1442
1443 self.event_publisher
1444 .publish(Event::RunStatusChanged(RunStatusChangedEvent {
1445 run_id,
1446 workflow_name: wf.clone(),
1447 from: RunStatus::Running,
1448 to,
1449 error: error.clone(),
1450 cost_usd,
1451 duration_ms,
1452 labels: labels.clone(),
1453 at: now,
1454 }));
1455
1456 if to == RunStatus::Failed {
1457 self.event_publisher
1458 .publish(Event::RunFailed(RunFailedEvent {
1459 run_id,
1460 workflow_name: wf,
1461 error,
1462 cost_usd,
1463 duration_ms,
1464 labels,
1465 at: now,
1466 }));
1467 }
1468 }
1469}
1470
1471impl fmt::Debug for Engine {
1472 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1473 f.debug_struct("Engine")
1474 .field("handlers", &self.handlers.keys().collect::<Vec<_>>())
1475 .finish_non_exhaustive()
1476 }
1477}
1478
1479#[cfg(test)]
1480mod tests {
1481 use super::*;
1482 use crate::config::ShellConfig;
1483 use crate::handler::{HandlerFuture, WorkflowHandler};
1484 use ironflow_core::providers::claude::ClaudeCodeProvider;
1485 use ironflow_core::providers::record_replay::RecordReplayProvider;
1486 use ironflow_store::memory::InMemoryStore;
1487 use ironflow_store::models::StepStatus;
1488 use serde_json::json;
1489
1490 struct EchoWorkflow;
1492
1493 impl WorkflowHandler for EchoWorkflow {
1494 fn name(&self) -> &str {
1495 "echo-workflow"
1496 }
1497
1498 fn describe(&self) -> WorkflowInfo {
1499 WorkflowInfo {
1500 description: "A simple workflow that echoes hello".to_string(),
1501 source_code: None,
1502 sub_workflows: Vec::new(),
1503 category: None,
1504 version: self.version().map(str::to_string),
1505 compatible_versions: Vec::new(),
1506 input_schema: None,
1507 default_labels: HashMap::new(),
1508 schedule: self.schedule().cloned(),
1509 default_max_cost_usd: self.default_max_cost_usd(),
1510 }
1511 }
1512
1513 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
1514 Box::pin(async move {
1515 ctx.shell("greet", ShellConfig::new("echo hello")).await?;
1516 Ok(())
1517 })
1518 }
1519 }
1520
1521 struct FailingWorkflow;
1523
1524 impl WorkflowHandler for FailingWorkflow {
1525 fn name(&self) -> &str {
1526 "failing-workflow"
1527 }
1528
1529 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
1530 Box::pin(async move {
1531 ctx.shell("fail", ShellConfig::new("exit 1")).await?;
1532 Ok(())
1533 })
1534 }
1535 }
1536
1537 fn create_test_engine() -> Engine {
1538 let store = Arc::new(InMemoryStore::new());
1539 let inner = ClaudeCodeProvider::new();
1540 let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
1541 inner,
1542 "/tmp/ironflow-fixtures",
1543 ));
1544 Engine::new(store, provider)
1545 }
1546
1547 #[test]
1548 fn engine_new_creates_instance() {
1549 let engine = create_test_engine();
1550 assert_eq!(engine.handler_names().len(), 0);
1551 }
1552
1553 #[test]
1554 fn engine_register_handler() {
1555 let mut engine = create_test_engine();
1556 let result = engine.register(EchoWorkflow);
1557 assert!(result.is_ok());
1558 assert_eq!(engine.handler_names().len(), 1);
1559 assert!(engine.handler_names().contains(&"echo-workflow"));
1560 }
1561
1562 #[test]
1563 fn engine_register_duplicate_returns_error() {
1564 let mut engine = create_test_engine();
1565 engine.register(EchoWorkflow).unwrap();
1566 let result = engine.register(EchoWorkflow);
1567 assert!(result.is_err());
1568 }
1569
1570 #[test]
1571 fn engine_get_handler_found() {
1572 let mut engine = create_test_engine();
1573 engine.register(EchoWorkflow).unwrap();
1574 let handler = engine.get_handler("echo-workflow");
1575 assert!(handler.is_some());
1576 }
1577
1578 #[test]
1579 fn engine_get_handler_not_found() {
1580 let engine = create_test_engine();
1581 let handler = engine.get_handler("nonexistent");
1582 assert!(handler.is_none());
1583 }
1584
1585 #[test]
1586 fn engine_handler_names_lists_all() {
1587 let mut engine = create_test_engine();
1588 engine.register(EchoWorkflow).unwrap();
1589 engine.register(FailingWorkflow).unwrap();
1590 let names = engine.handler_names();
1591 assert_eq!(names.len(), 2);
1592 assert!(names.contains(&"echo-workflow"));
1593 assert!(names.contains(&"failing-workflow"));
1594 }
1595
1596 #[test]
1597 fn engine_handler_info_returns_description() {
1598 let mut engine = create_test_engine();
1599 engine.register(EchoWorkflow).unwrap();
1600 let info = engine.handler_info("echo-workflow");
1601 assert!(info.is_some());
1602 let info = info.unwrap();
1603 assert_eq!(info.description, "A simple workflow that echoes hello");
1604 }
1605
1606 struct CategorizedWorkflow;
1607
1608 impl WorkflowHandler for CategorizedWorkflow {
1609 fn name(&self) -> &str {
1610 "categorized"
1611 }
1612 fn category(&self) -> Option<&str> {
1613 Some("data/etl")
1614 }
1615 fn execute<'a>(
1616 &'a self,
1617 _ctx: &'a mut WorkflowContext,
1618 ) -> crate::handler::HandlerFuture<'a> {
1619 Box::pin(async move { Ok(()) })
1620 }
1621 }
1622
1623 #[test]
1624 fn engine_default_describe_propagates_category() {
1625 let mut engine = create_test_engine();
1626 engine.register(CategorizedWorkflow).unwrap();
1627 let info = engine.handler_info("categorized").unwrap();
1628 assert_eq!(info.category.as_deref(), Some("data/etl"));
1629 }
1630
1631 #[test]
1632 fn engine_default_describe_without_category() {
1633 let mut engine = create_test_engine();
1634 engine.register(EchoWorkflow).unwrap();
1635 let info = engine.handler_info("echo-workflow").unwrap();
1636 assert!(info.category.is_none());
1637 }
1638
1639 struct ScheduledWorkflow {
1644 schedule: CronSchedule,
1645 }
1646
1647 impl ScheduledWorkflow {
1648 fn new() -> Self {
1649 Self {
1650 schedule: CronSchedule::new("0 0 * * * *").unwrap(),
1651 }
1652 }
1653 }
1654
1655 impl WorkflowHandler for ScheduledWorkflow {
1656 fn name(&self) -> &str {
1657 "scheduled"
1658 }
1659 fn schedule(&self) -> Option<&CronSchedule> {
1660 Some(&self.schedule)
1661 }
1662 fn execute<'a>(
1663 &'a self,
1664 _ctx: &'a mut WorkflowContext,
1665 ) -> crate::handler::HandlerFuture<'a> {
1666 Box::pin(async move { Ok(()) })
1667 }
1668 }
1669
1670 #[test]
1671 fn engine_default_describe_propagates_schedule() {
1672 let mut engine = create_test_engine();
1673 engine.register(ScheduledWorkflow::new()).unwrap();
1674 let info = engine.handler_info("scheduled").unwrap();
1675 assert_eq!(
1676 info.schedule.as_ref().map(|s| s.as_str()),
1677 Some("0 0 * * * *")
1678 );
1679 }
1680
1681 #[test]
1682 fn engine_default_describe_without_schedule() {
1683 let mut engine = create_test_engine();
1684 engine.register(EchoWorkflow).unwrap();
1685 let info = engine.handler_info("echo-workflow").unwrap();
1686 assert!(info.schedule.is_none());
1687 }
1688
1689 #[test]
1690 fn scheduled_handlers_returns_only_scheduled() {
1691 let mut engine = create_test_engine();
1692 engine.register(EchoWorkflow).unwrap();
1693 engine.register(ScheduledWorkflow::new()).unwrap();
1694 engine.register(FailingWorkflow).unwrap();
1695
1696 let scheduled = engine.scheduled_handlers();
1697 assert_eq!(scheduled.len(), 1);
1698 assert_eq!(scheduled[0].0, "scheduled");
1699 assert_eq!(scheduled[0].1.as_str(), "0 0 * * * *");
1700 }
1701
1702 #[test]
1703 fn scheduled_handlers_empty_when_none_scheduled() {
1704 let mut engine = create_test_engine();
1705 engine.register(EchoWorkflow).unwrap();
1706 engine.register(FailingWorkflow).unwrap();
1707
1708 let scheduled = engine.scheduled_handlers();
1709 assert!(scheduled.is_empty());
1710 }
1711
1712 struct BadCategoryWorkflow(&'static str);
1713
1714 impl WorkflowHandler for BadCategoryWorkflow {
1715 fn name(&self) -> &str {
1716 "bad-category"
1717 }
1718 fn category(&self) -> Option<&str> {
1719 Some(self.0)
1720 }
1721 fn execute<'a>(
1722 &'a self,
1723 _ctx: &'a mut WorkflowContext,
1724 ) -> crate::handler::HandlerFuture<'a> {
1725 Box::pin(async move { Ok(()) })
1726 }
1727 }
1728
1729 #[test]
1730 fn engine_register_rejects_empty_category() {
1731 let mut engine = create_test_engine();
1732 let err = engine.register(BadCategoryWorkflow("")).unwrap_err();
1733 match err {
1734 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty category")),
1735 other => panic!("expected InvalidWorkflow, got {other:?}"),
1736 }
1737 }
1738
1739 #[test]
1740 fn engine_register_rejects_leading_slash_category() {
1741 let mut engine = create_test_engine();
1742 let err = engine
1743 .register(BadCategoryWorkflow("/data/etl"))
1744 .unwrap_err();
1745 match err {
1746 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("leading '/'")),
1747 other => panic!("expected InvalidWorkflow, got {other:?}"),
1748 }
1749 }
1750
1751 #[test]
1752 fn engine_register_rejects_trailing_slash_category() {
1753 let mut engine = create_test_engine();
1754 let err = engine
1755 .register(BadCategoryWorkflow("data/etl/"))
1756 .unwrap_err();
1757 match err {
1758 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("trailing '/'")),
1759 other => panic!("expected InvalidWorkflow, got {other:?}"),
1760 }
1761 }
1762
1763 #[test]
1764 fn engine_register_rejects_double_slash_category() {
1765 let mut engine = create_test_engine();
1766 let err = engine
1767 .register(BadCategoryWorkflow("data//etl"))
1768 .unwrap_err();
1769 match err {
1770 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty segment")),
1771 other => panic!("expected InvalidWorkflow, got {other:?}"),
1772 }
1773 }
1774
1775 #[test]
1776 fn engine_register_rejects_whitespace_only_segment_category() {
1777 let mut engine = create_test_engine();
1778 let err = engine
1779 .register(BadCategoryWorkflow("data/ /etl"))
1780 .unwrap_err();
1781 match err {
1782 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("whitespace-only segment")),
1783 other => panic!("expected InvalidWorkflow, got {other:?}"),
1784 }
1785 }
1786
1787 #[test]
1788 fn engine_register_accepts_valid_nested_category() {
1789 let mut engine = create_test_engine();
1790 assert!(engine.register(CategorizedWorkflow).is_ok());
1791 }
1792
1793 #[tokio::test]
1794 async fn engine_unknown_workflow_returns_error() {
1795 let engine = create_test_engine();
1796 let result = engine
1797 .run_handler("unknown", TriggerKind::Manual, json!({}))
1798 .await;
1799 assert!(result.is_err());
1800 match result {
1801 Err(EngineError::InvalidWorkflow(msg)) => {
1802 assert!(msg.contains("no handler registered"));
1803 }
1804 _ => panic!("expected InvalidWorkflow error"),
1805 }
1806 }
1807
1808 #[tokio::test]
1809 async fn engine_enqueue_handler_creates_pending_run() {
1810 let mut engine = create_test_engine();
1811 engine.register(EchoWorkflow).unwrap();
1812
1813 let run = engine
1814 .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
1815 .await
1816 .unwrap();
1817 assert_eq!(run.status.state, RunStatus::Pending);
1818 assert_eq!(run.workflow_name, "echo-workflow");
1819 }
1820
1821 #[tokio::test]
1822 async fn enqueue_handler_leaves_the_run_unattributed() {
1823 let mut engine = create_test_engine();
1824 engine.register(EchoWorkflow).unwrap();
1825
1826 let run = engine
1827 .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
1828 .await
1829 .unwrap();
1830
1831 assert!(run.created_by.is_none());
1832 }
1833
1834 #[tokio::test]
1835 async fn enqueue_handler_with_options_records_the_author() {
1836 let mut engine = create_test_engine();
1837 engine.register(EchoWorkflow).unwrap();
1838 let actor = RunActor::User {
1839 user_id: Uuid::now_v7(),
1840 };
1841
1842 let run = engine
1843 .enqueue_handler_with_options(
1844 "echo-workflow",
1845 TriggerKind::Api,
1846 json!({}),
1847 EnqueueOptions {
1848 created_by: Some(actor.clone()),
1849 ..Default::default()
1850 },
1851 )
1852 .await
1853 .unwrap()
1854 .into_run();
1855
1856 assert_eq!(run.created_by, Some(actor));
1857 }
1858
1859 #[tokio::test]
1860 async fn enqueue_handler_with_options_accepts_no_author() {
1861 let mut engine = create_test_engine();
1862 engine.register(EchoWorkflow).unwrap();
1863
1864 let run = engine
1865 .enqueue_handler_with_options(
1866 "echo-workflow",
1867 TriggerKind::Cron {
1868 schedule: "0 * * * * *".to_string(),
1869 },
1870 json!({}),
1871 EnqueueOptions::default(),
1872 )
1873 .await
1874 .unwrap()
1875 .into_run();
1876
1877 assert!(run.created_by.is_none());
1878 }
1879
1880 #[tokio::test]
1881 async fn run_handler_leaves_the_run_unattributed() {
1882 let mut engine = create_test_engine();
1883 engine.register(EchoWorkflow).unwrap();
1884
1885 let run = engine
1886 .run_handler("echo-workflow", TriggerKind::Manual, json!({}))
1887 .await
1888 .unwrap()
1889 .run;
1890
1891 assert!(run.created_by.is_none());
1892 }
1893
1894 #[tokio::test]
1895 async fn engine_register_boxed() {
1896 let mut engine = create_test_engine();
1897 let handler: Box<dyn WorkflowHandler> = Box::new(EchoWorkflow);
1898 let result = engine.register_boxed(handler);
1899 assert!(result.is_ok());
1900 assert_eq!(engine.handler_names().len(), 1);
1901 }
1902
1903 #[tokio::test]
1904 async fn engine_store_and_provider_accessors() {
1905 let store = Arc::new(InMemoryStore::new());
1906 let inner = ClaudeCodeProvider::new();
1907 let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
1908 inner,
1909 "/tmp/ironflow-fixtures",
1910 ));
1911 let engine = Engine::new(store.clone(), provider.clone());
1912
1913 let _ = engine.store();
1915 let _ = engine.provider();
1916 }
1917
1918 use crate::operation::{Operation, OperationContext};
1923 use async_trait::async_trait;
1924 use ironflow_core::error::OperationError;
1925 use ironflow_store::models::StepKind;
1926
1927 struct FakeGitlabOp {
1928 project_id: u64,
1929 title: String,
1930 }
1931
1932 #[async_trait]
1933 impl Operation for FakeGitlabOp {
1934 fn kind(&self) -> &str {
1935 "gitlab"
1936 }
1937
1938 async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
1939 Ok(json!({
1940 "issue_id": 42,
1941 "project_id": self.project_id,
1942 "title": self.title,
1943 }))
1944 }
1945
1946 fn input(&self) -> Option<Value> {
1947 Some(json!({
1948 "project_id": self.project_id,
1949 "title": self.title,
1950 }))
1951 }
1952 }
1953
1954 struct FailingOp;
1955
1956 #[async_trait]
1957 impl Operation for FailingOp {
1958 fn kind(&self) -> &str {
1959 "broken-service"
1960 }
1961
1962 async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
1963 Err(OperationError::Http {
1964 status: None,
1965 message: "service unavailable".to_string(),
1966 })
1967 }
1968 }
1969
1970 struct OperationWorkflow;
1971
1972 impl WorkflowHandler for OperationWorkflow {
1973 fn name(&self) -> &str {
1974 "operation-workflow"
1975 }
1976
1977 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
1978 Box::pin(async move {
1979 let op = FakeGitlabOp {
1980 project_id: 123,
1981 title: "Bug report".to_string(),
1982 };
1983 ctx.operation("create-issue", &op).await?;
1984 Ok(())
1985 })
1986 }
1987 }
1988
1989 struct FailingOperationWorkflow;
1990
1991 impl WorkflowHandler for FailingOperationWorkflow {
1992 fn name(&self) -> &str {
1993 "failing-operation-workflow"
1994 }
1995
1996 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
1997 Box::pin(async move {
1998 ctx.operation("broken-call", &FailingOp).await?;
1999 Ok(())
2000 })
2001 }
2002 }
2003
2004 struct MixedWorkflow;
2005
2006 impl WorkflowHandler for MixedWorkflow {
2007 fn name(&self) -> &str {
2008 "mixed-workflow"
2009 }
2010
2011 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2012 Box::pin(async move {
2013 ctx.shell("build", ShellConfig::new("echo built")).await?;
2014 let op = FakeGitlabOp {
2015 project_id: 456,
2016 title: "Deploy done".to_string(),
2017 };
2018 let result = ctx.operation("notify-gitlab", &op).await?;
2019 assert_eq!(result.output["issue_id"], 42);
2020 Ok(())
2021 })
2022 }
2023 }
2024
2025 #[tokio::test]
2026 async fn operation_step_happy_path() {
2027 let mut engine = create_test_engine();
2028 engine.register(OperationWorkflow).unwrap();
2029
2030 let run = engine
2031 .run_handler("operation-workflow", TriggerKind::Manual, json!({}))
2032 .await
2033 .unwrap()
2034 .run;
2035
2036 assert_eq!(run.status.state, RunStatus::Completed);
2037
2038 let steps = engine.store().list_steps(run.id).await.unwrap();
2039
2040 assert_eq!(steps.len(), 1);
2041 assert_eq!(steps[0].name, "create-issue");
2042 assert_eq!(steps[0].kind, StepKind::Custom("gitlab".to_string()));
2043 assert_eq!(
2044 steps[0].status.state,
2045 ironflow_store::models::StepStatus::Completed
2046 );
2047
2048 let output = steps[0].output.as_ref().unwrap();
2049 assert_eq!(output["issue_id"], 42);
2050 assert_eq!(output["project_id"], 123);
2051
2052 let input = steps[0].input.as_ref().unwrap();
2053 assert_eq!(input["project_id"], 123);
2054 assert_eq!(input["title"], "Bug report");
2055 }
2056
2057 #[tokio::test]
2058 async fn operation_step_failure_marks_run_failed() {
2059 let mut engine = create_test_engine();
2060 engine.register(FailingOperationWorkflow).unwrap();
2061
2062 let result = engine
2063 .run_handler("failing-operation-workflow", TriggerKind::Manual, json!({}))
2064 .await;
2065
2066 assert!(result.is_err());
2067 }
2068
2069 #[tokio::test]
2070 async fn operation_mixed_with_shell_steps() {
2071 let mut engine = create_test_engine();
2072 engine.register(MixedWorkflow).unwrap();
2073
2074 let run = engine
2075 .run_handler("mixed-workflow", TriggerKind::Manual, json!({}))
2076 .await
2077 .unwrap()
2078 .run;
2079
2080 assert_eq!(run.status.state, RunStatus::Completed);
2081
2082 let steps = engine.store().list_steps(run.id).await.unwrap();
2083
2084 assert_eq!(steps.len(), 2);
2085 assert_eq!(steps[0].kind, StepKind::Shell);
2086 assert_eq!(steps[1].kind, StepKind::Custom("gitlab".to_string()));
2087 assert_eq!(steps[0].position, 0);
2088 assert_eq!(steps[1].position, 1);
2089 }
2090
2091 use crate::config::ApprovalConfig;
2096
2097 struct SingleApprovalWorkflow;
2098
2099 impl WorkflowHandler for SingleApprovalWorkflow {
2100 fn name(&self) -> &str {
2101 "single-approval"
2102 }
2103
2104 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2105 Box::pin(async move {
2106 ctx.shell("build", ShellConfig::new("echo built")).await?;
2107 ctx.approval("gate", ApprovalConfig::new("OK?")).await?;
2108 ctx.shell("deploy", ShellConfig::new("echo deployed"))
2109 .await?;
2110 Ok(())
2111 })
2112 }
2113 }
2114
2115 struct DoubleApprovalWorkflow;
2116
2117 impl WorkflowHandler for DoubleApprovalWorkflow {
2118 fn name(&self) -> &str {
2119 "double-approval"
2120 }
2121
2122 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2123 Box::pin(async move {
2124 ctx.shell("build", ShellConfig::new("echo built")).await?;
2125 ctx.approval("staging-gate", ApprovalConfig::new("Deploy staging?"))
2126 .await?;
2127 ctx.shell("deploy-staging", ShellConfig::new("echo staging"))
2128 .await?;
2129 ctx.approval("prod-gate", ApprovalConfig::new("Deploy prod?"))
2130 .await?;
2131 ctx.shell("deploy-prod", ShellConfig::new("echo prod"))
2132 .await?;
2133 Ok(())
2134 })
2135 }
2136 }
2137
2138 #[tokio::test]
2139 async fn approval_pauses_run() {
2140 let mut engine = create_test_engine();
2141 engine.register(SingleApprovalWorkflow).unwrap();
2142
2143 let run = engine
2144 .run_handler("single-approval", TriggerKind::Manual, json!({}))
2145 .await
2146 .unwrap()
2147 .run;
2148
2149 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
2150
2151 let steps = engine.store().list_steps(run.id).await.unwrap();
2152 assert_eq!(steps.len(), 2); assert_eq!(steps[0].kind, StepKind::Shell);
2154 assert_eq!(steps[0].status.state, StepStatus::Completed);
2155 assert_eq!(steps[1].kind, StepKind::Approval);
2156 assert_eq!(steps[1].status.state, StepStatus::AwaitingApproval);
2157 }
2158
2159 #[tokio::test]
2160 async fn approval_resume_completes_run() {
2161 let mut engine = create_test_engine();
2162 engine.register(SingleApprovalWorkflow).unwrap();
2163
2164 let run = engine
2166 .run_handler("single-approval", TriggerKind::Manual, json!({}))
2167 .await
2168 .unwrap()
2169 .run;
2170 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
2171
2172 engine
2174 .store()
2175 .update_run_status(run.id, RunStatus::Running)
2176 .await
2177 .unwrap();
2178
2179 let resumed = engine.resume_run(run.id).await.unwrap().run;
2181 assert_eq!(resumed.status.state, RunStatus::Completed);
2182
2183 let steps = engine.store().list_steps(run.id).await.unwrap();
2184 assert_eq!(steps.len(), 3); assert_eq!(steps[0].name, "build");
2186 assert_eq!(steps[0].status.state, StepStatus::Completed);
2187 assert_eq!(steps[1].name, "gate");
2188 assert_eq!(steps[1].kind, StepKind::Approval);
2189 assert_eq!(steps[1].status.state, StepStatus::Completed);
2190 assert_eq!(steps[2].name, "deploy");
2191 assert_eq!(steps[2].status.state, StepStatus::Completed);
2192 }
2193
2194 #[tokio::test]
2195 async fn double_approval_two_resumes() {
2196 let mut engine = create_test_engine();
2197 engine.register(DoubleApprovalWorkflow).unwrap();
2198
2199 let run = engine
2201 .run_handler("double-approval", TriggerKind::Manual, json!({}))
2202 .await
2203 .unwrap()
2204 .run;
2205 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
2206
2207 let steps = engine.store().list_steps(run.id).await.unwrap();
2208 assert_eq!(steps.len(), 2); engine
2212 .store()
2213 .update_run_status(run.id, RunStatus::Running)
2214 .await
2215 .unwrap();
2216
2217 let resumed = engine.resume_run(run.id).await.unwrap().run;
2218 assert_eq!(resumed.status.state, RunStatus::AwaitingApproval);
2219
2220 let steps = engine.store().list_steps(run.id).await.unwrap();
2221 assert_eq!(steps.len(), 4); engine
2225 .store()
2226 .update_run_status(run.id, RunStatus::Running)
2227 .await
2228 .unwrap();
2229
2230 let final_run = engine.resume_run(run.id).await.unwrap().run;
2231 assert_eq!(final_run.status.state, RunStatus::Completed);
2232
2233 let steps = engine.store().list_steps(run.id).await.unwrap();
2234 assert_eq!(steps.len(), 5);
2235 assert_eq!(steps[0].name, "build");
2236 assert_eq!(steps[1].name, "staging-gate");
2237 assert_eq!(steps[2].name, "deploy-staging");
2238 assert_eq!(steps[3].name, "prod-gate");
2239 assert_eq!(steps[4].name, "deploy-prod");
2240
2241 for step in &steps {
2242 assert_eq!(step.status.state, StepStatus::Completed);
2243 }
2244 }
2245
2246 use ironflow_store::models::{NewStep, StepUpdate, step_trace_id};
2251
2252 async fn create_step_with_status(
2253 store: &Arc<dyn Store>,
2254 run_id: Uuid,
2255 name: &str,
2256 position: u32,
2257 status: StepStatus,
2258 ) -> ironflow_store::models::Step {
2259 let step = store
2260 .create_step(NewStep {
2261 run_id,
2262 trace_id: step_trace_id(run_id, name, position),
2263 name: name.to_string(),
2264 kind: StepKind::Shell,
2265 position,
2266 input: None,
2267 is_error_handler: false,
2268 })
2269 .await
2270 .unwrap();
2271
2272 match status {
2273 StepStatus::Pending => {}
2274 StepStatus::Running => {
2275 store
2276 .update_step(
2277 step.id,
2278 StepUpdate {
2279 status: Some(StepStatus::Running),
2280 ..StepUpdate::default()
2281 },
2282 )
2283 .await
2284 .unwrap();
2285 }
2286 StepStatus::Completed => {
2287 store
2288 .update_step(
2289 step.id,
2290 StepUpdate {
2291 status: Some(StepStatus::Running),
2292 ..StepUpdate::default()
2293 },
2294 )
2295 .await
2296 .unwrap();
2297 store
2298 .update_step(
2299 step.id,
2300 StepUpdate {
2301 status: Some(StepStatus::Completed),
2302 ..StepUpdate::default()
2303 },
2304 )
2305 .await
2306 .unwrap();
2307 }
2308 StepStatus::AwaitingApproval => {
2309 store
2310 .update_step(
2311 step.id,
2312 StepUpdate {
2313 status: Some(StepStatus::Running),
2314 ..StepUpdate::default()
2315 },
2316 )
2317 .await
2318 .unwrap();
2319 store
2320 .update_step(
2321 step.id,
2322 StepUpdate {
2323 status: Some(StepStatus::AwaitingApproval),
2324 ..StepUpdate::default()
2325 },
2326 )
2327 .await
2328 .unwrap();
2329 }
2330 _ => panic!("unsupported status for test helper: {status}"),
2331 }
2332
2333 store.get_step(step.id).await.unwrap().unwrap()
2334 }
2335
2336 #[tokio::test]
2337 async fn fail_orphaned_steps_marks_running_as_failed() {
2338 let engine = create_test_engine();
2339 let run = engine
2340 .store()
2341 .create_run(NewRun {
2342 created_by: None,
2343 workflow_name: "test".to_string(),
2344 trigger: TriggerKind::Manual,
2345 payload: json!({}),
2346 max_retries: 0,
2347 handler_version: None,
2348 labels: HashMap::new(),
2349 scheduled_at: None,
2350 idempotency_key: None,
2351 max_cost_usd: None,
2352 })
2353 .await
2354 .unwrap()
2355 .into_run();
2356
2357 let step = create_step_with_status(
2358 engine.store(),
2359 run.id,
2360 "running-step",
2361 0,
2362 StepStatus::Running,
2363 )
2364 .await;
2365
2366 engine
2367 .fail_orphaned_steps(run.id, "parent run timed out")
2368 .await
2369 .unwrap();
2370
2371 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
2372 assert_eq!(updated.status.state, StepStatus::Failed);
2373 assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
2374 assert!(updated.completed_at.is_some());
2375 }
2376
2377 #[tokio::test]
2378 async fn fail_orphaned_steps_marks_pending_as_skipped() {
2379 let engine = create_test_engine();
2380 let run = engine
2381 .store()
2382 .create_run(NewRun {
2383 created_by: None,
2384 workflow_name: "test".to_string(),
2385 trigger: TriggerKind::Manual,
2386 payload: json!({}),
2387 max_retries: 0,
2388 handler_version: None,
2389 labels: HashMap::new(),
2390 scheduled_at: None,
2391 idempotency_key: None,
2392 max_cost_usd: None,
2393 })
2394 .await
2395 .unwrap()
2396 .into_run();
2397
2398 let step = create_step_with_status(
2399 engine.store(),
2400 run.id,
2401 "pending-step",
2402 0,
2403 StepStatus::Pending,
2404 )
2405 .await;
2406
2407 engine
2408 .fail_orphaned_steps(run.id, "parent run timed out")
2409 .await
2410 .unwrap();
2411
2412 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
2413 assert_eq!(updated.status.state, StepStatus::Skipped);
2414 assert!(updated.error.is_none());
2415 assert!(updated.completed_at.is_some());
2416 }
2417
2418 #[tokio::test]
2419 async fn fail_orphaned_steps_marks_awaiting_approval_as_failed() {
2420 let engine = create_test_engine();
2421 let run = engine
2422 .store()
2423 .create_run(NewRun {
2424 created_by: None,
2425 workflow_name: "test".to_string(),
2426 trigger: TriggerKind::Manual,
2427 payload: json!({}),
2428 max_retries: 0,
2429 handler_version: None,
2430 labels: HashMap::new(),
2431 scheduled_at: None,
2432 idempotency_key: None,
2433 max_cost_usd: None,
2434 })
2435 .await
2436 .unwrap()
2437 .into_run();
2438
2439 let step = create_step_with_status(
2440 engine.store(),
2441 run.id,
2442 "approval-step",
2443 0,
2444 StepStatus::AwaitingApproval,
2445 )
2446 .await;
2447
2448 engine
2449 .fail_orphaned_steps(run.id, "parent run timed out")
2450 .await
2451 .unwrap();
2452
2453 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
2454 assert_eq!(updated.status.state, StepStatus::Failed);
2455 assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
2456 assert!(updated.completed_at.is_some());
2457 }
2458
2459 #[tokio::test]
2460 async fn fail_orphaned_steps_skips_terminal_steps() {
2461 let engine = create_test_engine();
2462 let run = engine
2463 .store()
2464 .create_run(NewRun {
2465 created_by: None,
2466 workflow_name: "test".to_string(),
2467 trigger: TriggerKind::Manual,
2468 payload: json!({}),
2469 max_retries: 0,
2470 handler_version: None,
2471 labels: HashMap::new(),
2472 scheduled_at: None,
2473 idempotency_key: None,
2474 max_cost_usd: None,
2475 })
2476 .await
2477 .unwrap()
2478 .into_run();
2479
2480 let completed_step =
2481 create_step_with_status(engine.store(), run.id, "done", 0, StepStatus::Completed).await;
2482 let running_step =
2483 create_step_with_status(engine.store(), run.id, "in-flight", 1, StepStatus::Running)
2484 .await;
2485
2486 engine
2487 .fail_orphaned_steps(run.id, "parent run timed out")
2488 .await
2489 .unwrap();
2490
2491 let completed = engine
2492 .store()
2493 .get_step(completed_step.id)
2494 .await
2495 .unwrap()
2496 .unwrap();
2497 assert_eq!(completed.status.state, StepStatus::Completed);
2498
2499 let failed = engine
2500 .store()
2501 .get_step(running_step.id)
2502 .await
2503 .unwrap()
2504 .unwrap();
2505 assert_eq!(failed.status.state, StepStatus::Failed);
2506 }
2507
2508 #[tokio::test]
2509 async fn fail_orphaned_steps_mixed_states() {
2510 let engine = create_test_engine();
2511 let run = engine
2512 .store()
2513 .create_run(NewRun {
2514 created_by: None,
2515 workflow_name: "test".to_string(),
2516 trigger: TriggerKind::Manual,
2517 payload: json!({}),
2518 max_retries: 0,
2519 handler_version: None,
2520 labels: HashMap::new(),
2521 scheduled_at: None,
2522 idempotency_key: None,
2523 max_cost_usd: None,
2524 })
2525 .await
2526 .unwrap()
2527 .into_run();
2528
2529 let s_completed =
2530 create_step_with_status(engine.store(), run.id, "step-1", 0, StepStatus::Completed)
2531 .await;
2532 let s_running =
2533 create_step_with_status(engine.store(), run.id, "step-2", 1, StepStatus::Running).await;
2534 let s_pending =
2535 create_step_with_status(engine.store(), run.id, "step-3", 2, StepStatus::Pending).await;
2536
2537 engine.fail_orphaned_steps(run.id, "timeout").await.unwrap();
2538
2539 let r_completed = engine
2540 .store()
2541 .get_step(s_completed.id)
2542 .await
2543 .unwrap()
2544 .unwrap();
2545 assert_eq!(r_completed.status.state, StepStatus::Completed);
2546
2547 let r_running = engine
2548 .store()
2549 .get_step(s_running.id)
2550 .await
2551 .unwrap()
2552 .unwrap();
2553 assert_eq!(r_running.status.state, StepStatus::Failed);
2554 assert_eq!(r_running.error.as_deref(), Some("timeout"));
2555
2556 let r_pending = engine
2557 .store()
2558 .get_step(s_pending.id)
2559 .await
2560 .unwrap()
2561 .unwrap();
2562 assert_eq!(r_pending.status.state, StepStatus::Skipped);
2563 assert!(r_pending.error.is_none());
2564 }
2565
2566 #[tokio::test]
2567 async fn fail_orphaned_steps_no_steps_is_noop() {
2568 let engine = create_test_engine();
2569 let run = engine
2570 .store()
2571 .create_run(NewRun {
2572 created_by: None,
2573 workflow_name: "test".to_string(),
2574 trigger: TriggerKind::Manual,
2575 payload: json!({}),
2576 max_retries: 0,
2577 handler_version: None,
2578 labels: HashMap::new(),
2579 scheduled_at: None,
2580 idempotency_key: None,
2581 max_cost_usd: None,
2582 })
2583 .await
2584 .unwrap()
2585 .into_run();
2586
2587 let result = engine.fail_orphaned_steps(run.id, "timeout").await;
2588 assert!(result.is_ok());
2589 }
2590
2591 #[tokio::test]
2592 async fn fail_orphaned_steps_preserves_existing_error() {
2593 let engine = create_test_engine();
2594 let run = engine
2595 .store()
2596 .create_run(NewRun {
2597 created_by: None,
2598 workflow_name: "test".to_string(),
2599 trigger: TriggerKind::Manual,
2600 payload: json!({}),
2601 max_retries: 0,
2602 handler_version: None,
2603 labels: HashMap::new(),
2604 scheduled_at: None,
2605 idempotency_key: None,
2606 max_cost_usd: None,
2607 })
2608 .await
2609 .unwrap()
2610 .into_run();
2611
2612 let step_with_error = create_step_with_status(
2613 engine.store(),
2614 run.id,
2615 "already-errored",
2616 0,
2617 StepStatus::Running,
2618 )
2619 .await;
2620
2621 engine
2622 .store()
2623 .update_step(
2624 step_with_error.id,
2625 StepUpdate {
2626 error: Some("real error from provider".to_string()),
2627 ..StepUpdate::default()
2628 },
2629 )
2630 .await
2631 .unwrap();
2632
2633 let step_no_error = create_step_with_status(
2634 engine.store(),
2635 run.id,
2636 "no-error-yet",
2637 1,
2638 StepStatus::Running,
2639 )
2640 .await;
2641
2642 engine
2643 .fail_orphaned_steps(run.id, "parent run failed")
2644 .await
2645 .unwrap();
2646
2647 let updated_with = engine
2648 .store()
2649 .get_step(step_with_error.id)
2650 .await
2651 .unwrap()
2652 .unwrap();
2653 assert_eq!(updated_with.status.state, StepStatus::Failed);
2654 assert_eq!(
2655 updated_with.error.as_deref(),
2656 Some("real error from provider"),
2657 );
2658
2659 let updated_without = engine
2660 .store()
2661 .get_step(step_no_error.id)
2662 .await
2663 .unwrap()
2664 .unwrap();
2665 assert_eq!(updated_without.status.state, StepStatus::Failed);
2666 assert_eq!(updated_without.error.as_deref(), Some("parent run failed"),);
2667 }
2668}