1use std::collections::HashMap;
10use std::fmt;
11use std::sync::{Arc, Mutex};
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
20use ironflow_core::error::OperationError;
21#[cfg(feature = "prometheus")]
22use ironflow_core::metric_names::{
23 RUN_BUDGET_EXCEEDED_TOTAL, RUN_COST_USD, RUN_DURATION_SECONDS, RUNS_ACTIVE, RUNS_TOTAL,
24};
25use ironflow_core::provider::AgentProvider;
26use ironflow_store::error::StoreError;
27use ironflow_store::models::{
28 NewRun, Run, RunActor, RunCreation, RunFilter, RunStatus, RunUpdate, StepStatus, StepUpdate,
29 TriggerKind,
30};
31use ironflow_store::store::Store;
32#[cfg(feature = "prometheus")]
33use metrics::{counter, gauge, histogram};
34
35use crate::artifact::ArtifactSink;
36use crate::budget::{BudgetConfig, month_start};
37use crate::context::WorkflowContext;
38use crate::error::EngineError;
39use crate::executor::{StepInterceptor, StepResult};
40use crate::guard::{WorkflowGuardConfig, new_shared_guard_state};
41use crate::handler::{WorkflowHandler, WorkflowInfo};
42use crate::log_sender::LogSender;
43use crate::notify::{
44 ApprovalRequestedEvent, Event, EventPublisher, EventSubscriber, RunBudgetExceededEvent,
45 RunFailedEvent, RunStatusChangedEvent, WorkflowEventBus,
46};
47use crate::plan::{
48 ExecutionPlan, PlanOptions, PlanRecorder, SharedPlanRecorder, estimate_durations, lock_plan,
49};
50use crate::retry_policy::{backoff_for_retry, is_run_retryable};
51use crate::schedule::CronSchedule;
52use ironflow_core::decision::DecisionProvider;
53
54#[derive(Debug, Clone)]
73pub struct WorkflowResult {
74 pub run: Run,
76 pub steps: Vec<StepResult>,
78}
79
80#[derive(Debug, Clone, Default)]
99pub struct EnqueueOptions {
100 pub max_retries: u32,
102 pub labels: HashMap<String, String>,
104 pub scheduled_at: Option<DateTime<Utc>>,
107 pub max_cost_usd: Option<Decimal>,
111 pub created_by: Option<RunActor>,
114 pub idempotency_key: Option<String>,
120}
121
122#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
133pub enum ExecutionMode {
134 #[default]
138 Local,
139 Workers,
142}
143
144pub struct Engine {
185 store: Arc<dyn Store>,
186 provider: Arc<dyn AgentProvider>,
187 handlers: HashMap<String, Arc<dyn WorkflowHandler>>,
188 event_publisher: EventPublisher,
189 log_sender: Option<LogSender>,
190 budget: BudgetConfig,
191 artifact_sink: Option<Arc<dyn ArtifactSink>>,
192 guard_config: Option<WorkflowGuardConfig>,
193 event_bus: Option<WorkflowEventBus>,
194 decision_provider: Option<Arc<dyn DecisionProvider>>,
195 step_interceptor: Option<Arc<dyn StepInterceptor>>,
196 execution_mode: ExecutionMode,
197}
198
199fn validate_category(handler_name: &str, category: &str) -> Result<(), EngineError> {
209 let reject = |reason: &str| {
210 Err(EngineError::InvalidWorkflow(format!(
211 "handler '{handler_name}' has invalid category '{category}': {reason}"
212 )))
213 };
214
215 if category.is_empty() {
216 return reject("empty category");
217 }
218 if category.starts_with('/') {
219 return reject("leading '/'");
220 }
221 if category.ends_with('/') {
222 return reject("trailing '/'");
223 }
224 for segment in category.split('/') {
225 if segment.is_empty() {
226 return reject("empty segment (double '/')");
227 }
228 if segment.trim().is_empty() {
229 return reject("whitespace-only segment");
230 }
231 }
232 Ok(())
233}
234
235impl Engine {
236 pub fn new(store: Arc<dyn Store>, provider: Arc<dyn AgentProvider>) -> Self {
252 Self {
253 store,
254 provider,
255 handlers: HashMap::new(),
256 event_publisher: EventPublisher::new(),
257 log_sender: None,
258 budget: BudgetConfig::new(),
259 artifact_sink: None,
260 guard_config: None,
261 event_bus: None,
262 decision_provider: None,
263 step_interceptor: None,
264 execution_mode: ExecutionMode::default(),
265 }
266 }
267
268 pub fn with_decision_provider(mut self, provider: Arc<dyn DecisionProvider>) -> Self {
289 self.decision_provider = Some(provider);
290 self
291 }
292
293 pub fn with_step_interceptor(mut self, interceptor: Arc<dyn StepInterceptor>) -> Self {
317 self.step_interceptor = Some(interceptor);
318 self
319 }
320
321 pub fn step_interceptor(&self) -> Option<&Arc<dyn StepInterceptor>> {
323 self.step_interceptor.as_ref()
324 }
325
326 pub fn with_budget_config(mut self, budget: BudgetConfig) -> Self {
347 self.budget = budget;
348 self
349 }
350
351 pub fn budget_config(&self) -> &BudgetConfig {
353 &self.budget
354 }
355
356 pub fn with_guard_config(mut self, config: WorkflowGuardConfig) -> Self {
378 self.guard_config = Some(config);
379 self
380 }
381
382 pub fn guard_config(&self) -> Option<&WorkflowGuardConfig> {
384 self.guard_config.as_ref()
385 }
386
387 pub fn with_execution_mode(mut self, mode: ExecutionMode) -> Self {
408 self.execution_mode = mode;
409 self
410 }
411
412 pub fn execution_mode(&self) -> ExecutionMode {
414 self.execution_mode
415 }
416
417 pub fn set_log_sender(&mut self, sender: LogSender) {
423 self.log_sender = Some(sender);
424 }
425
426 pub fn set_artifact_sink(&mut self, sink: Arc<dyn ArtifactSink>) {
444 self.artifact_sink = Some(sink);
445 }
446
447 pub fn artifact_sink(&self) -> Option<&Arc<dyn ArtifactSink>> {
449 self.artifact_sink.as_ref()
450 }
451
452 pub fn set_event_bus(&mut self, bus: WorkflowEventBus) {
469 self.event_bus = Some(bus);
470 }
471
472 pub fn event_bus(&self) -> Option<&WorkflowEventBus> {
474 self.event_bus.as_ref()
475 }
476
477 pub fn store(&self) -> &Arc<dyn Store> {
479 &self.store
480 }
481
482 pub fn provider(&self) -> &Arc<dyn AgentProvider> {
484 &self.provider
485 }
486
487 fn build_context(&self, run: &Run) -> WorkflowContext {
496 let handlers = self.handlers.clone();
497 let resolver: crate::context::HandlerResolver =
498 Arc::new(move |name: &str| handlers.get(name).cloned());
499 let mut ctx = WorkflowContext::with_handler_resolver(
500 run.id,
501 run.workflow_name.clone(),
502 self.store.clone(),
503 self.provider.clone(),
504 resolver,
505 );
506 ctx.carry_over_run_totals(run.retry_count + 1, run.cost_usd, run.duration_ms);
507 ctx.set_max_cost_usd(run.max_cost_usd);
508 if let Some(ref sender) = self.log_sender {
509 ctx.set_log_sender(sender.clone());
510 }
511 if let Some(ref sink) = self.artifact_sink {
512 ctx.set_artifact_sink(sink.clone());
513 }
514 if let Some(ref bus) = self.event_bus {
515 ctx.set_event_bus(bus.clone());
516 }
517 if let Some(ref provider) = self.decision_provider {
518 ctx.set_decision_provider(provider.clone());
519 }
520 if let Some(ref interceptor) = self.step_interceptor {
521 ctx.set_step_interceptor(interceptor.clone());
522 }
523 ctx
524 }
525
526 fn build_context_with_guard(
532 &self,
533 run: &Run,
534 handler: &dyn WorkflowHandler,
535 ) -> WorkflowContext {
536 let mut ctx = self.build_context(run);
537 let guard_config = handler.guard_config().or_else(|| self.guard_config.clone());
538 if let Some(config) = guard_config {
539 ctx.set_guard(config, new_shared_guard_state());
540 }
541 ctx
542 }
543
544 async fn check_monthly_quota(&self, workflow_name: &str) -> Result<(), EngineError> {
555 let Some(limit) = self.budget.monthly_cost_limit_usd else {
556 return Ok(());
557 };
558
559 let stats = self
560 .store
561 .get_stats(RunFilter {
562 created_after: Some(month_start(Utc::now())),
563 ..RunFilter::default()
564 })
565 .await?;
566
567 if stats.total_cost_usd < limit {
568 return Ok(());
569 }
570
571 warn!(
572 workflow = %workflow_name,
573 limit_usd = %limit,
574 spent_usd = %stats.total_cost_usd,
575 "monthly cost quota exhausted, refusing new run"
576 );
577
578 #[cfg(feature = "prometheus")]
579 counter!(
580 RUN_BUDGET_EXCEEDED_TOTAL,
581 "workflow" => workflow_name.to_string(),
582 "scope" => "monthly",
583 )
584 .increment(1);
585
586 Err(EngineError::MonthlyBudgetExceeded {
587 limit_usd: limit,
588 spent_usd: stats.total_cost_usd,
589 })
590 }
591
592 pub fn register(&mut self, handler: impl WorkflowHandler + 'static) -> Result<(), EngineError> {
636 let name = handler.name().to_string();
637 if self.handlers.contains_key(&name) {
638 return Err(EngineError::InvalidWorkflow(format!(
639 "handler '{}' already registered",
640 name
641 )));
642 }
643 if let Some(category) = handler.category() {
644 validate_category(&name, category)?;
645 }
646 self.handlers.insert(name, Arc::new(handler));
647 Ok(())
648 }
649
650 pub fn register_boxed(&mut self, handler: Box<dyn WorkflowHandler>) -> Result<(), EngineError> {
657 let name = handler.name().to_string();
658 if self.handlers.contains_key(&name) {
659 return Err(EngineError::InvalidWorkflow(format!(
660 "handler '{}' already registered",
661 name
662 )));
663 }
664 if let Some(category) = handler.category() {
665 validate_category(&name, category)?;
666 }
667 self.handlers.insert(name, Arc::from(handler));
668 Ok(())
669 }
670
671 pub fn get_handler(&self, name: &str) -> Option<&Arc<dyn WorkflowHandler>> {
673 self.handlers.get(name)
674 }
675
676 pub fn handler_names(&self) -> Vec<&str> {
678 self.handlers.keys().map(|s| s.as_str()).collect()
679 }
680
681 pub fn handler_info(&self, name: &str) -> Option<WorkflowInfo> {
683 self.handlers.get(name).map(|h| h.describe())
684 }
685
686 pub fn scheduled_handlers(&self) -> Vec<(&str, &CronSchedule)> {
710 self.handlers
711 .iter()
712 .filter_map(|(name, handler)| handler.schedule().map(|sched| (name.as_str(), sched)))
713 .collect()
714 }
715
716 pub fn subscribe(
741 &mut self,
742 subscriber: impl EventSubscriber + 'static,
743 event_types: &[&'static str],
744 ) {
745 self.event_publisher.subscribe(subscriber, event_types);
746 }
747
748 pub fn event_publisher(&self) -> &EventPublisher {
753 &self.event_publisher
754 }
755
756 #[tracing::instrument(name = "engine.run_handler", skip_all, fields(workflow = %handler_name))]
786 pub async fn run_handler(
787 &self,
788 handler_name: &str,
789 trigger: TriggerKind,
790 payload: Value,
791 ) -> Result<WorkflowResult, EngineError> {
792 let handler = self
793 .handlers
794 .get(handler_name)
795 .ok_or_else(|| {
796 EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
797 })?
798 .clone();
799
800 self.check_monthly_quota(handler_name).await?;
801
802 let handler_version = handler.version().map(str::to_string);
803 let max_cost_usd = self
804 .budget
805 .resolve_run_cap(None, handler.default_max_cost_usd());
806 let run = self
807 .store
808 .create_run(NewRun {
809 created_by: None,
810 workflow_name: handler_name.to_string(),
811 trigger,
812 payload,
813 max_retries: 0,
814 handler_version,
815 labels: handler.default_labels(),
816 scheduled_at: None,
817 idempotency_key: None,
818 max_cost_usd,
819 })
820 .await?
821 .into_run();
822
823 let run_id = run.id;
824 info!(run_id = %run_id, handler_version = run.handler_version.as_deref().unwrap_or(""), "run created");
825
826 self.store
827 .update_run_status(run_id, RunStatus::Running)
828 .await?;
829
830 #[cfg(feature = "prometheus")]
831 gauge!(RUNS_ACTIVE, "workflow" => handler_name.to_string()).increment(1.0);
832
833 let run_start = Instant::now();
834 let mut ctx = self.build_context_with_guard(&run, handler.as_ref());
835
836 let result = handler.execute(&mut ctx).await;
837 self.finalize_run(run_id, handler_name, result, &ctx, run_start, run.labels)
838 .await
839 }
840
841 #[tracing::instrument(name = "engine.plan_handler", skip_all, fields(workflow = %handler_name))]
883 pub async fn plan_handler(
884 &self,
885 handler_name: &str,
886 payload: Value,
887 options: PlanOptions,
888 ) -> Result<ExecutionPlan, EngineError> {
889 if options.max_depth == 0 {
890 return Err(EngineError::InvalidWorkflow(
891 "max_depth must be at least 1".to_string(),
892 ));
893 }
894
895 let handler = self
896 .handlers
897 .get(handler_name)
898 .ok_or_else(|| {
899 EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
900 })?
901 .clone();
902
903 let estimates = if options.estimate_durations {
904 estimate_durations(&self.store, handler_name, options.sample_runs).await?
905 } else {
906 HashMap::new()
907 };
908
909 let shared: SharedPlanRecorder = Arc::new(Mutex::new(PlanRecorder::new(
910 handler_name.to_string(),
911 payload,
912 options.max_depth,
913 estimates,
914 )));
915
916 let handlers = self.handlers.clone();
919 let resolver: crate::context::HandlerResolver =
920 Arc::new(move |name: &str| handlers.get(name).cloned());
921 let mut ctx = WorkflowContext::with_handler_resolver(
922 Uuid::now_v7(),
923 handler_name.to_string(),
924 self.store.clone(),
925 self.provider.clone(),
926 resolver,
927 );
928 ctx.set_plan(shared.clone());
929
930 if let Err(err) = handler.execute(&mut ctx).await {
931 lock_plan(&shared).fail(err.to_string());
932 }
933 drop(ctx);
934
935 let plan = match Arc::try_unwrap(shared) {
936 Ok(mutex) => mutex
937 .into_inner()
938 .unwrap_or_else(|poisoned| poisoned.into_inner())
939 .into_plan(),
940 Err(shared) => lock_plan(&shared).snapshot(),
941 };
942
943 info!(
944 workflow = %handler_name,
945 steps = plan.steps.len(),
946 truncated = plan.truncated,
947 "execution plan built"
948 );
949
950 Ok(plan)
951 }
952
953 #[tracing::instrument(name = "engine.enqueue_handler", skip_all, fields(workflow = %handler_name))]
964 pub async fn enqueue_handler(
965 &self,
966 handler_name: &str,
967 trigger: TriggerKind,
968 payload: Value,
969 max_retries: u32,
970 ) -> Result<Run, EngineError> {
971 self.enqueue_handler_with_options(
972 handler_name,
973 trigger,
974 payload,
975 EnqueueOptions {
976 max_retries,
977 ..Default::default()
978 },
979 )
980 .await
981 .map(RunCreation::into_run)
982 }
983
984 #[tracing::instrument(name = "engine.enqueue_handler_with_options", skip_all, fields(workflow = %handler_name))]
1029 pub async fn enqueue_handler_with_options(
1030 &self,
1031 handler_name: &str,
1032 trigger: TriggerKind,
1033 payload: Value,
1034 options: EnqueueOptions,
1035 ) -> Result<RunCreation, EngineError> {
1036 let EnqueueOptions {
1037 max_retries,
1038 labels,
1039 scheduled_at,
1040 max_cost_usd,
1041 created_by,
1042 idempotency_key,
1043 } = options;
1044
1045 let handler = self.handlers.get(handler_name).ok_or_else(|| {
1046 EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
1047 })?;
1048
1049 self.check_monthly_quota(handler_name).await?;
1050
1051 let handler_version = handler.version().map(str::to_string);
1052 let mut merged_labels = handler.default_labels();
1053 merged_labels.extend(labels);
1054 let resolved_cap = self
1055 .budget
1056 .resolve_run_cap(max_cost_usd, handler.default_max_cost_usd());
1057
1058 let creation = self
1059 .store
1060 .create_run(NewRun {
1061 workflow_name: handler_name.to_string(),
1062 trigger,
1063 payload,
1064 max_retries,
1065 handler_version,
1066 labels: merged_labels,
1067 scheduled_at,
1068 created_by,
1069 idempotency_key,
1070 max_cost_usd: resolved_cap,
1071 })
1072 .await?;
1073
1074 match &creation {
1075 RunCreation::Created(run) => info!(
1076 run_id = %run.id,
1077 workflow = %handler_name,
1078 max_cost_usd = ?resolved_cap,
1079 "handler run enqueued"
1080 ),
1081 RunCreation::Existing(run) => info!(
1082 run_id = %run.id,
1083 workflow = %handler_name,
1084 "idempotent replay, nothing enqueued"
1085 ),
1086 }
1087
1088 Ok(creation)
1089 }
1090
1091 #[tracing::instrument(name = "engine.execute_handler_run", skip_all, fields(run_id = %run_id))]
1104 pub async fn execute_handler_run(&self, run_id: Uuid) -> Result<WorkflowResult, EngineError> {
1105 let run = self
1106 .store
1107 .get_run(run_id)
1108 .await?
1109 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
1110
1111 let handler = self
1112 .handlers
1113 .get(&run.workflow_name)
1114 .ok_or_else(|| {
1115 EngineError::InvalidWorkflow(format!(
1116 "no handler registered: {}",
1117 run.workflow_name
1118 ))
1119 })?
1120 .clone();
1121
1122 #[cfg(feature = "prometheus")]
1123 gauge!(RUNS_ACTIVE, "workflow" => run.workflow_name.clone()).increment(1.0);
1124
1125 let run_start = Instant::now();
1126 let mut ctx = self.build_context_with_guard(&run, handler.as_ref());
1127
1128 ctx.load_replay_steps().await?;
1134
1135 let result = self
1136 .release_then_execute(run_id, handler.as_ref(), &mut ctx)
1137 .await;
1138 self.finalize_run(
1139 run_id,
1140 &run.workflow_name,
1141 result,
1142 &ctx,
1143 run_start,
1144 run.labels,
1145 )
1146 .await
1147 }
1148
1149 #[tracing::instrument(name = "engine.execute_run", skip_all, fields(run_id = %run_id))]
1157 pub async fn execute_run(&self, run_id: Uuid) -> Result<WorkflowResult, EngineError> {
1158 self.execute_handler_run(run_id).await
1159 }
1160
1161 #[tracing::instrument(name = "engine.resume_run", skip_all, fields(run_id = %run_id))]
1180 pub async fn resume_run(&self, run_id: Uuid) -> Result<WorkflowResult, EngineError> {
1181 let run = self
1182 .store
1183 .get_run(run_id)
1184 .await?
1185 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
1186
1187 let handler = self
1188 .handlers
1189 .get(&run.workflow_name)
1190 .ok_or_else(|| {
1191 EngineError::InvalidWorkflow(format!(
1192 "no handler registered: {}",
1193 run.workflow_name
1194 ))
1195 })?
1196 .clone();
1197
1198 info!(run_id = %run_id, workflow = %run.workflow_name, "resuming run after approval");
1199
1200 let run_start = Instant::now();
1201 let mut ctx = self.build_context_with_guard(&run, handler.as_ref());
1202 ctx.load_replay_steps().await?;
1203
1204 let result = self
1205 .release_then_execute(run_id, handler.as_ref(), &mut ctx)
1206 .await;
1207 self.finalize_run(
1208 run_id,
1209 &run.workflow_name,
1210 result,
1211 &ctx,
1212 run_start,
1213 run.labels,
1214 )
1215 .await
1216 }
1217
1218 pub async fn fail_or_schedule_retry(
1263 &self,
1264 run_id: Uuid,
1265 error: &str,
1266 retryable: bool,
1267 cost_usd: Option<Decimal>,
1268 duration_ms: Option<u64>,
1269 ) -> Result<RunStatus, EngineError> {
1270 let run = self
1271 .store
1272 .get_run(run_id)
1273 .await?
1274 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
1275
1276 let has_attempts_left = run.retry_count < run.max_retries;
1277 let update = if retryable && has_attempts_left {
1278 let backoff = backoff_for_retry(run.retry_count);
1279 let scheduled_at = Utc::now() + TimeDelta::milliseconds(backoff.as_millis() as i64);
1280
1281 info!(
1282 run_id = %run_id,
1283 workflow = %run.workflow_name,
1284 attempt = run.retry_count + 1,
1285 max_retries = run.max_retries,
1286 backoff_secs = backoff.as_secs(),
1287 scheduled_at = %scheduled_at,
1288 "run failed, scheduling retry"
1289 );
1290
1291 RunUpdate {
1292 status: Some(RunStatus::Retrying),
1293 error: Some(error.to_string()),
1294 increment_retry: true,
1295 cost_usd,
1296 duration_ms,
1297 scheduled_at: Some(scheduled_at),
1298 ..RunUpdate::default()
1299 }
1300 } else {
1301 RunUpdate {
1302 status: Some(RunStatus::Failed),
1303 error: Some(error.to_string()),
1304 cost_usd,
1305 duration_ms,
1306 completed_at: Some(Utc::now()),
1307 ..RunUpdate::default()
1308 }
1309 };
1310
1311 let status = update.status.unwrap_or(RunStatus::Failed);
1312 self.store.update_run(run_id, update).await?;
1313 self.fail_orphaned_steps(run_id, error).await?;
1314
1315 Ok(status)
1316 }
1317
1318 pub async fn fail_orphaned_steps(
1332 &self,
1333 run_id: Uuid,
1334 error_message: &str,
1335 ) -> Result<(), EngineError> {
1336 let steps = self.store.list_steps(run_id).await?;
1337 let now = Utc::now();
1338
1339 for step in steps {
1340 if step.status.state.is_terminal() {
1341 continue;
1342 }
1343
1344 let (target_status, error) = match step.status.state {
1345 StepStatus::Running | StepStatus::AwaitingApproval => {
1346 let err = if step.error.is_some() {
1347 None
1348 } else {
1349 Some(error_message.to_string())
1350 };
1351 (StepStatus::Failed, err)
1352 }
1353 StepStatus::Pending => (StepStatus::Skipped, None),
1354 _ => continue,
1355 };
1356
1357 if let Err(e) = self
1358 .store
1359 .update_step(
1360 step.id,
1361 StepUpdate {
1362 status: Some(target_status),
1363 error,
1364 completed_at: Some(now),
1365 ..StepUpdate::default()
1366 },
1367 )
1368 .await
1369 {
1370 warn!(
1371 run_id = %run_id,
1372 step_id = %step.id,
1373 step_name = %step.name,
1374 error = %e,
1375 "failed to cleanup orphaned step"
1376 );
1377 } else {
1378 info!(
1379 run_id = %run_id,
1380 step_id = %step.id,
1381 step_name = %step.name,
1382 from = %step.status.state,
1383 to = %target_status,
1384 "cleaned up orphaned step"
1385 );
1386 }
1387 }
1388
1389 Ok(())
1390 }
1391
1392 async fn release_then_execute(
1397 &self,
1398 run_id: Uuid,
1399 handler: &dyn WorkflowHandler,
1400 ctx: &mut WorkflowContext,
1401 ) -> Result<(), EngineError> {
1402 match self.provider.release_run(&run_id.to_string()).await {
1403 Ok(()) => handler.execute(ctx).await,
1404 Err(e) => Err(EngineError::Operation(OperationError::Agent(e))),
1405 }
1406 }
1407
1408 async fn finalize_run(
1414 &self,
1415 run_id: Uuid,
1416 workflow_name: &str,
1417 result: Result<(), EngineError>,
1418 ctx: &WorkflowContext,
1419 run_start: Instant,
1420 run_labels: HashMap<String, String>,
1421 ) -> Result<WorkflowResult, EngineError> {
1422 let total_duration = ctx.carried_duration_ms() + run_start.elapsed().as_millis() as u64;
1425 let completed_at = Utc::now();
1426
1427 let final_status;
1428 let final_run;
1429
1430 match result {
1431 Ok(()) => {
1432 final_status = if ctx.has_allowed_failure() {
1433 RunStatus::Warning
1434 } else {
1435 RunStatus::Completed
1436 };
1437 final_run = self
1438 .store
1439 .update_run_returning(
1440 run_id,
1441 RunUpdate {
1442 status: Some(final_status),
1443 cost_usd: Some(ctx.total_cost_usd()),
1444 duration_ms: Some(total_duration),
1445 completed_at: Some(completed_at),
1446 ..RunUpdate::default()
1447 },
1448 )
1449 .await?;
1450
1451 info!(
1452 run_id = %run_id,
1453 status = %final_status,
1454 cost_usd = %ctx.total_cost_usd(),
1455 duration_ms = total_duration,
1456 "run completed"
1457 );
1458 }
1459 Err(EngineError::ApprovalRequired {
1460 run_id: approval_run_id,
1461 step_id,
1462 ref message,
1463 }) => {
1464 final_status = RunStatus::AwaitingApproval;
1465 final_run = self
1466 .store
1467 .update_run_returning(
1468 run_id,
1469 RunUpdate {
1470 status: Some(RunStatus::AwaitingApproval),
1471 cost_usd: Some(ctx.total_cost_usd()),
1472 duration_ms: Some(total_duration),
1473 ..RunUpdate::default()
1474 },
1475 )
1476 .await?;
1477
1478 info!(
1479 run_id = %approval_run_id,
1480 step_id = %step_id,
1481 message = %message,
1482 "run awaiting approval"
1483 );
1484
1485 let requirement = self
1487 .store
1488 .get_step(step_id)
1489 .await?
1490 .and_then(|s| s.approval_requirement);
1491 self.event_publisher
1492 .publish(Event::ApprovalRequested(ApprovalRequestedEvent {
1493 run_id: approval_run_id,
1494 step_id,
1495 message: message.clone(),
1496 requirement,
1497 at: Utc::now(),
1498 }));
1499 }
1500 Err(EngineError::HumanInputRequired {
1501 run_id: input_run_id,
1502 step_id,
1503 ref message,
1504 }) => {
1505 final_status = RunStatus::AwaitingApproval;
1506 final_run = self
1507 .store
1508 .update_run_returning(
1509 run_id,
1510 RunUpdate {
1511 status: Some(RunStatus::AwaitingApproval),
1512 cost_usd: Some(ctx.total_cost_usd()),
1513 duration_ms: Some(total_duration),
1514 ..RunUpdate::default()
1515 },
1516 )
1517 .await?;
1518
1519 info!(
1521 run_id = %input_run_id,
1522 step_id = %step_id,
1523 message = %message,
1524 "run awaiting human input"
1525 );
1526 }
1527 Err(EngineError::DelaySleeping {
1528 run_id: delay_run_id,
1529 step_id,
1530 wake_at,
1531 }) => {
1532 final_status = RunStatus::Sleeping;
1533 final_run = self
1534 .store
1535 .update_run_returning(
1536 run_id,
1537 RunUpdate {
1538 status: Some(RunStatus::Sleeping),
1539 cost_usd: Some(ctx.total_cost_usd()),
1540 duration_ms: Some(total_duration),
1541 scheduled_at: Some(wake_at),
1542 ..RunUpdate::default()
1543 },
1544 )
1545 .await?;
1546
1547 info!(
1548 run_id = %delay_run_id,
1549 step_id = %step_id,
1550 wake_at = %wake_at,
1551 "run sleeping until delay elapses"
1552 );
1553 }
1554 Err(err) => {
1555 let guardrail_stop = matches!(
1559 err,
1560 EngineError::RunBudgetExceeded { .. } | EngineError::WorkflowGuardRejected(_)
1561 );
1562
1563 final_status = if guardrail_stop {
1564 if let Err(store_err) = self
1565 .store
1566 .update_run(
1567 run_id,
1568 RunUpdate {
1569 status: Some(RunStatus::Cancelled),
1570 error: Some(err.to_string()),
1571 cost_usd: Some(ctx.total_cost_usd()),
1572 duration_ms: Some(total_duration),
1573 completed_at: Some(completed_at),
1574 ..RunUpdate::default()
1575 },
1576 )
1577 .await
1578 {
1579 error!(run_id = %run_id, store_error = %store_err, "failed to persist run cancellation");
1580 }
1581 if let Err(cleanup_err) = self
1582 .fail_orphaned_steps(run_id, "run stopped: guardrail limit reached")
1583 .await
1584 {
1585 error!(run_id = %run_id, store_error = %cleanup_err, "failed to cleanup orphaned steps");
1586 }
1587 RunStatus::Cancelled
1588 } else {
1589 self.fail_or_schedule_retry(
1590 run_id,
1591 &err.to_string(),
1592 is_run_retryable(&err),
1593 Some(ctx.total_cost_usd()),
1594 Some(total_duration),
1595 )
1596 .await
1597 .unwrap_or_else(|store_err| {
1598 error!(run_id = %run_id, store_error = %store_err, "failed to persist run failure");
1599 RunStatus::Failed
1600 })
1601 };
1602
1603 if matches!(err, EngineError::RunBudgetExceeded { .. }) {
1604 self.on_run_budget_exceeded(workflow_name, run_id, &err);
1605 }
1606
1607 error!(run_id = %run_id, status = %final_status, error = %err, "run stopped");
1608
1609 self.publish_run_status_changed(
1610 workflow_name,
1611 run_id,
1612 final_status,
1613 Some(err.to_string()),
1614 ctx,
1615 total_duration,
1616 run_labels,
1617 );
1618
1619 #[cfg(feature = "prometheus")]
1620 self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
1621
1622 return Err(err);
1623 }
1624 }
1625
1626 self.publish_run_status_changed(
1627 workflow_name,
1628 run_id,
1629 final_status,
1630 None,
1631 ctx,
1632 total_duration,
1633 run_labels,
1634 );
1635
1636 #[cfg(feature = "prometheus")]
1637 self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
1638
1639 Ok(WorkflowResult {
1640 run: final_run,
1641 steps: ctx.step_results().to_vec(),
1642 })
1643 }
1644
1645 #[cfg(feature = "prometheus")]
1647 fn emit_run_metrics(
1648 &self,
1649 workflow_name: &str,
1650 status: RunStatus,
1651 duration_ms: u64,
1652 ctx: &WorkflowContext,
1653 ) {
1654 let status_str = status.to_string();
1655 let wf = workflow_name.to_string();
1656
1657 counter!(RUNS_TOTAL, "workflow" => wf.clone(), "status" => status_str.clone()).increment(1);
1658 histogram!(RUN_DURATION_SECONDS, "workflow" => wf.clone(), "status" => status_str)
1659 .record(duration_ms as f64 / 1000.0);
1660 histogram!(RUN_COST_USD, "workflow" => wf.clone()).record(
1661 ctx.total_cost_usd()
1662 .to_string()
1663 .parse::<f64>()
1664 .unwrap_or(0.0),
1665 );
1666 gauge!(RUNS_ACTIVE, "workflow" => wf).decrement(1.0);
1667 }
1668
1669 fn on_run_budget_exceeded(&self, workflow_name: &str, run_id: Uuid, err: &EngineError) {
1675 let EngineError::RunBudgetExceeded {
1676 limit_usd,
1677 spent_usd,
1678 step_budget_usd,
1679 ..
1680 } = err
1681 else {
1682 return;
1683 };
1684
1685 #[cfg(feature = "prometheus")]
1686 counter!(
1687 RUN_BUDGET_EXCEEDED_TOTAL,
1688 "workflow" => workflow_name.to_string(),
1689 "scope" => "run",
1690 )
1691 .increment(1);
1692
1693 self.event_publisher
1694 .publish(Event::RunBudgetExceeded(RunBudgetExceededEvent {
1695 run_id,
1696 workflow_name: workflow_name.to_string(),
1697 limit_usd: *limit_usd,
1698 spent_usd: *spent_usd,
1699 step_budget_usd: *step_budget_usd,
1700 at: Utc::now(),
1701 }));
1702 }
1703
1704 #[allow(clippy::too_many_arguments)]
1709 fn publish_run_status_changed(
1710 &self,
1711 workflow_name: &str,
1712 run_id: Uuid,
1713 to: RunStatus,
1714 error: Option<String>,
1715 ctx: &WorkflowContext,
1716 duration_ms: u64,
1717 labels: HashMap<String, String>,
1718 ) {
1719 let now = Utc::now();
1720 let cost_usd = ctx.total_cost_usd();
1721 let wf = workflow_name.to_string();
1722
1723 self.event_publisher
1724 .publish(Event::RunStatusChanged(RunStatusChangedEvent {
1725 run_id,
1726 workflow_name: wf.clone(),
1727 from: RunStatus::Running,
1728 to,
1729 error: error.clone(),
1730 cost_usd,
1731 duration_ms,
1732 labels: labels.clone(),
1733 at: now,
1734 }));
1735
1736 if to == RunStatus::Failed {
1737 self.event_publisher
1738 .publish(Event::RunFailed(RunFailedEvent {
1739 run_id,
1740 workflow_name: wf,
1741 error,
1742 cost_usd,
1743 duration_ms,
1744 labels,
1745 at: now,
1746 }));
1747 }
1748 }
1749}
1750
1751impl fmt::Debug for Engine {
1752 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1753 f.debug_struct("Engine")
1754 .field("handlers", &self.handlers.keys().collect::<Vec<_>>())
1755 .finish_non_exhaustive()
1756 }
1757}
1758
1759#[cfg(test)]
1760mod tests {
1761 use super::*;
1762 use crate::config::ShellConfig;
1763 use crate::handler::{HandlerFuture, WorkflowHandler};
1764 use ironflow_core::providers::claude::ClaudeCodeProvider;
1765 use ironflow_core::providers::record_replay::RecordReplayProvider;
1766 use ironflow_store::memory::InMemoryStore;
1767 use ironflow_store::models::StepStatus;
1768 use serde_json::json;
1769
1770 struct EchoWorkflow;
1772
1773 impl WorkflowHandler for EchoWorkflow {
1774 fn name(&self) -> &str {
1775 "echo-workflow"
1776 }
1777
1778 fn describe(&self) -> WorkflowInfo {
1779 WorkflowInfo {
1780 description: "A simple workflow that echoes hello".to_string(),
1781 source_code: None,
1782 sub_workflows: Vec::new(),
1783 category: None,
1784 version: self.version().map(str::to_string),
1785 compatible_versions: Vec::new(),
1786 input_schema: None,
1787 default_labels: HashMap::new(),
1788 schedule: self.schedule().cloned(),
1789 default_max_cost_usd: self.default_max_cost_usd(),
1790 }
1791 }
1792
1793 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
1794 Box::pin(async move {
1795 ctx.shell("greet", ShellConfig::new("echo hello")).await?;
1796 Ok(())
1797 })
1798 }
1799 }
1800
1801 struct FailingWorkflow;
1803
1804 impl WorkflowHandler for FailingWorkflow {
1805 fn name(&self) -> &str {
1806 "failing-workflow"
1807 }
1808
1809 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
1810 Box::pin(async move {
1811 ctx.shell("fail", ShellConfig::new("exit 1")).await?;
1812 Ok(())
1813 })
1814 }
1815 }
1816
1817 fn create_test_engine() -> Engine {
1818 let store = Arc::new(InMemoryStore::new());
1819 let inner = ClaudeCodeProvider::new();
1820 let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
1821 inner,
1822 "/tmp/ironflow-fixtures",
1823 ));
1824 Engine::new(store, provider)
1825 }
1826
1827 #[test]
1828 fn engine_new_creates_instance() {
1829 let engine = create_test_engine();
1830 assert_eq!(engine.handler_names().len(), 0);
1831 }
1832
1833 #[test]
1834 fn execution_mode_defaults_to_local() {
1835 let engine = create_test_engine();
1836 assert_eq!(engine.execution_mode(), ExecutionMode::Local);
1837 }
1838
1839 #[test]
1840 fn with_execution_mode_overrides_the_default() {
1841 let engine = create_test_engine().with_execution_mode(ExecutionMode::Workers);
1842 assert_eq!(engine.execution_mode(), ExecutionMode::Workers);
1843 }
1844
1845 #[test]
1846 fn engine_register_handler() {
1847 let mut engine = create_test_engine();
1848 let result = engine.register(EchoWorkflow);
1849 assert!(result.is_ok());
1850 assert_eq!(engine.handler_names().len(), 1);
1851 assert!(engine.handler_names().contains(&"echo-workflow"));
1852 }
1853
1854 #[test]
1855 fn engine_register_duplicate_returns_error() {
1856 let mut engine = create_test_engine();
1857 engine.register(EchoWorkflow).unwrap();
1858 let result = engine.register(EchoWorkflow);
1859 assert!(result.is_err());
1860 }
1861
1862 #[test]
1863 fn engine_get_handler_found() {
1864 let mut engine = create_test_engine();
1865 engine.register(EchoWorkflow).unwrap();
1866 let handler = engine.get_handler("echo-workflow");
1867 assert!(handler.is_some());
1868 }
1869
1870 #[test]
1871 fn engine_get_handler_not_found() {
1872 let engine = create_test_engine();
1873 let handler = engine.get_handler("nonexistent");
1874 assert!(handler.is_none());
1875 }
1876
1877 #[test]
1878 fn engine_handler_names_lists_all() {
1879 let mut engine = create_test_engine();
1880 engine.register(EchoWorkflow).unwrap();
1881 engine.register(FailingWorkflow).unwrap();
1882 let names = engine.handler_names();
1883 assert_eq!(names.len(), 2);
1884 assert!(names.contains(&"echo-workflow"));
1885 assert!(names.contains(&"failing-workflow"));
1886 }
1887
1888 #[test]
1889 fn engine_handler_info_returns_description() {
1890 let mut engine = create_test_engine();
1891 engine.register(EchoWorkflow).unwrap();
1892 let info = engine.handler_info("echo-workflow");
1893 assert!(info.is_some());
1894 let info = info.unwrap();
1895 assert_eq!(info.description, "A simple workflow that echoes hello");
1896 }
1897
1898 struct CategorizedWorkflow;
1899
1900 impl WorkflowHandler for CategorizedWorkflow {
1901 fn name(&self) -> &str {
1902 "categorized"
1903 }
1904 fn category(&self) -> Option<&str> {
1905 Some("data/etl")
1906 }
1907 fn execute<'a>(
1908 &'a self,
1909 _ctx: &'a mut WorkflowContext,
1910 ) -> crate::handler::HandlerFuture<'a> {
1911 Box::pin(async move { Ok(()) })
1912 }
1913 }
1914
1915 #[test]
1916 fn engine_default_describe_propagates_category() {
1917 let mut engine = create_test_engine();
1918 engine.register(CategorizedWorkflow).unwrap();
1919 let info = engine.handler_info("categorized").unwrap();
1920 assert_eq!(info.category.as_deref(), Some("data/etl"));
1921 }
1922
1923 #[test]
1924 fn engine_default_describe_without_category() {
1925 let mut engine = create_test_engine();
1926 engine.register(EchoWorkflow).unwrap();
1927 let info = engine.handler_info("echo-workflow").unwrap();
1928 assert!(info.category.is_none());
1929 }
1930
1931 struct ScheduledWorkflow {
1936 schedule: CronSchedule,
1937 }
1938
1939 impl ScheduledWorkflow {
1940 fn new() -> Self {
1941 Self {
1942 schedule: CronSchedule::new("0 0 * * * *").unwrap(),
1943 }
1944 }
1945 }
1946
1947 impl WorkflowHandler for ScheduledWorkflow {
1948 fn name(&self) -> &str {
1949 "scheduled"
1950 }
1951 fn schedule(&self) -> Option<&CronSchedule> {
1952 Some(&self.schedule)
1953 }
1954 fn execute<'a>(
1955 &'a self,
1956 _ctx: &'a mut WorkflowContext,
1957 ) -> crate::handler::HandlerFuture<'a> {
1958 Box::pin(async move { Ok(()) })
1959 }
1960 }
1961
1962 #[test]
1963 fn engine_default_describe_propagates_schedule() {
1964 let mut engine = create_test_engine();
1965 engine.register(ScheduledWorkflow::new()).unwrap();
1966 let info = engine.handler_info("scheduled").unwrap();
1967 assert_eq!(
1968 info.schedule.as_ref().map(|s| s.as_str()),
1969 Some("0 0 * * * *")
1970 );
1971 }
1972
1973 #[test]
1974 fn engine_default_describe_without_schedule() {
1975 let mut engine = create_test_engine();
1976 engine.register(EchoWorkflow).unwrap();
1977 let info = engine.handler_info("echo-workflow").unwrap();
1978 assert!(info.schedule.is_none());
1979 }
1980
1981 #[test]
1982 fn scheduled_handlers_returns_only_scheduled() {
1983 let mut engine = create_test_engine();
1984 engine.register(EchoWorkflow).unwrap();
1985 engine.register(ScheduledWorkflow::new()).unwrap();
1986 engine.register(FailingWorkflow).unwrap();
1987
1988 let scheduled = engine.scheduled_handlers();
1989 assert_eq!(scheduled.len(), 1);
1990 assert_eq!(scheduled[0].0, "scheduled");
1991 assert_eq!(scheduled[0].1.as_str(), "0 0 * * * *");
1992 }
1993
1994 #[test]
1995 fn scheduled_handlers_empty_when_none_scheduled() {
1996 let mut engine = create_test_engine();
1997 engine.register(EchoWorkflow).unwrap();
1998 engine.register(FailingWorkflow).unwrap();
1999
2000 let scheduled = engine.scheduled_handlers();
2001 assert!(scheduled.is_empty());
2002 }
2003
2004 struct BadCategoryWorkflow(&'static str);
2005
2006 impl WorkflowHandler for BadCategoryWorkflow {
2007 fn name(&self) -> &str {
2008 "bad-category"
2009 }
2010 fn category(&self) -> Option<&str> {
2011 Some(self.0)
2012 }
2013 fn execute<'a>(
2014 &'a self,
2015 _ctx: &'a mut WorkflowContext,
2016 ) -> crate::handler::HandlerFuture<'a> {
2017 Box::pin(async move { Ok(()) })
2018 }
2019 }
2020
2021 #[test]
2022 fn engine_register_rejects_empty_category() {
2023 let mut engine = create_test_engine();
2024 let err = engine.register(BadCategoryWorkflow("")).unwrap_err();
2025 match err {
2026 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty category")),
2027 other => panic!("expected InvalidWorkflow, got {other:?}"),
2028 }
2029 }
2030
2031 #[test]
2032 fn engine_register_rejects_leading_slash_category() {
2033 let mut engine = create_test_engine();
2034 let err = engine
2035 .register(BadCategoryWorkflow("/data/etl"))
2036 .unwrap_err();
2037 match err {
2038 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("leading '/'")),
2039 other => panic!("expected InvalidWorkflow, got {other:?}"),
2040 }
2041 }
2042
2043 #[test]
2044 fn engine_register_rejects_trailing_slash_category() {
2045 let mut engine = create_test_engine();
2046 let err = engine
2047 .register(BadCategoryWorkflow("data/etl/"))
2048 .unwrap_err();
2049 match err {
2050 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("trailing '/'")),
2051 other => panic!("expected InvalidWorkflow, got {other:?}"),
2052 }
2053 }
2054
2055 #[test]
2056 fn engine_register_rejects_double_slash_category() {
2057 let mut engine = create_test_engine();
2058 let err = engine
2059 .register(BadCategoryWorkflow("data//etl"))
2060 .unwrap_err();
2061 match err {
2062 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty segment")),
2063 other => panic!("expected InvalidWorkflow, got {other:?}"),
2064 }
2065 }
2066
2067 #[test]
2068 fn engine_register_rejects_whitespace_only_segment_category() {
2069 let mut engine = create_test_engine();
2070 let err = engine
2071 .register(BadCategoryWorkflow("data/ /etl"))
2072 .unwrap_err();
2073 match err {
2074 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("whitespace-only segment")),
2075 other => panic!("expected InvalidWorkflow, got {other:?}"),
2076 }
2077 }
2078
2079 #[test]
2080 fn engine_register_accepts_valid_nested_category() {
2081 let mut engine = create_test_engine();
2082 assert!(engine.register(CategorizedWorkflow).is_ok());
2083 }
2084
2085 #[tokio::test]
2086 async fn engine_unknown_workflow_returns_error() {
2087 let engine = create_test_engine();
2088 let result = engine
2089 .run_handler("unknown", TriggerKind::Manual, json!({}))
2090 .await;
2091 assert!(result.is_err());
2092 match result {
2093 Err(EngineError::InvalidWorkflow(msg)) => {
2094 assert!(msg.contains("no handler registered"));
2095 }
2096 _ => panic!("expected InvalidWorkflow error"),
2097 }
2098 }
2099
2100 #[tokio::test]
2101 async fn engine_enqueue_handler_creates_pending_run() {
2102 let mut engine = create_test_engine();
2103 engine.register(EchoWorkflow).unwrap();
2104
2105 let run = engine
2106 .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
2107 .await
2108 .unwrap();
2109 assert_eq!(run.status.state, RunStatus::Pending);
2110 assert_eq!(run.workflow_name, "echo-workflow");
2111 }
2112
2113 #[tokio::test]
2114 async fn enqueue_handler_leaves_the_run_unattributed() {
2115 let mut engine = create_test_engine();
2116 engine.register(EchoWorkflow).unwrap();
2117
2118 let run = engine
2119 .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
2120 .await
2121 .unwrap();
2122
2123 assert!(run.created_by.is_none());
2124 }
2125
2126 #[tokio::test]
2127 async fn enqueue_handler_with_options_records_the_author() {
2128 let mut engine = create_test_engine();
2129 engine.register(EchoWorkflow).unwrap();
2130 let actor = RunActor::User {
2131 user_id: Uuid::now_v7(),
2132 };
2133
2134 let run = engine
2135 .enqueue_handler_with_options(
2136 "echo-workflow",
2137 TriggerKind::Api,
2138 json!({}),
2139 EnqueueOptions {
2140 created_by: Some(actor.clone()),
2141 ..Default::default()
2142 },
2143 )
2144 .await
2145 .unwrap()
2146 .into_run();
2147
2148 assert_eq!(run.created_by, Some(actor));
2149 }
2150
2151 #[tokio::test]
2152 async fn enqueue_handler_with_options_accepts_no_author() {
2153 let mut engine = create_test_engine();
2154 engine.register(EchoWorkflow).unwrap();
2155
2156 let run = engine
2157 .enqueue_handler_with_options(
2158 "echo-workflow",
2159 TriggerKind::Cron {
2160 schedule: "0 * * * * *".to_string(),
2161 },
2162 json!({}),
2163 EnqueueOptions::default(),
2164 )
2165 .await
2166 .unwrap()
2167 .into_run();
2168
2169 assert!(run.created_by.is_none());
2170 }
2171
2172 #[tokio::test]
2173 async fn run_handler_leaves_the_run_unattributed() {
2174 let mut engine = create_test_engine();
2175 engine.register(EchoWorkflow).unwrap();
2176
2177 let run = engine
2178 .run_handler("echo-workflow", TriggerKind::Manual, json!({}))
2179 .await
2180 .unwrap()
2181 .run;
2182
2183 assert!(run.created_by.is_none());
2184 }
2185
2186 #[tokio::test]
2187 async fn engine_register_boxed() {
2188 let mut engine = create_test_engine();
2189 let handler: Box<dyn WorkflowHandler> = Box::new(EchoWorkflow);
2190 let result = engine.register_boxed(handler);
2191 assert!(result.is_ok());
2192 assert_eq!(engine.handler_names().len(), 1);
2193 }
2194
2195 #[tokio::test]
2196 async fn engine_store_and_provider_accessors() {
2197 let store = Arc::new(InMemoryStore::new());
2198 let inner = ClaudeCodeProvider::new();
2199 let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
2200 inner,
2201 "/tmp/ironflow-fixtures",
2202 ));
2203 let engine = Engine::new(store.clone(), provider.clone());
2204
2205 let _ = engine.store();
2207 let _ = engine.provider();
2208 }
2209
2210 use crate::operation::{Operation, OperationContext};
2215 use async_trait::async_trait;
2216 use ironflow_core::error::OperationError;
2217 use ironflow_store::models::StepKind;
2218
2219 struct FakeGitlabOp {
2220 project_id: u64,
2221 title: String,
2222 }
2223
2224 #[async_trait]
2225 impl Operation for FakeGitlabOp {
2226 fn kind(&self) -> &str {
2227 "gitlab"
2228 }
2229
2230 async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
2231 Ok(json!({
2232 "issue_id": 42,
2233 "project_id": self.project_id,
2234 "title": self.title,
2235 }))
2236 }
2237
2238 fn input(&self) -> Option<Value> {
2239 Some(json!({
2240 "project_id": self.project_id,
2241 "title": self.title,
2242 }))
2243 }
2244 }
2245
2246 struct FailingOp;
2247
2248 #[async_trait]
2249 impl Operation for FailingOp {
2250 fn kind(&self) -> &str {
2251 "broken-service"
2252 }
2253
2254 async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
2255 Err(OperationError::Http {
2256 status: None,
2257 message: "service unavailable".to_string(),
2258 })
2259 }
2260 }
2261
2262 struct OperationWorkflow;
2263
2264 impl WorkflowHandler for OperationWorkflow {
2265 fn name(&self) -> &str {
2266 "operation-workflow"
2267 }
2268
2269 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2270 Box::pin(async move {
2271 let op = FakeGitlabOp {
2272 project_id: 123,
2273 title: "Bug report".to_string(),
2274 };
2275 ctx.operation("create-issue", &op).await?;
2276 Ok(())
2277 })
2278 }
2279 }
2280
2281 struct FailingOperationWorkflow;
2282
2283 impl WorkflowHandler for FailingOperationWorkflow {
2284 fn name(&self) -> &str {
2285 "failing-operation-workflow"
2286 }
2287
2288 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2289 Box::pin(async move {
2290 ctx.operation("broken-call", &FailingOp).await?;
2291 Ok(())
2292 })
2293 }
2294 }
2295
2296 struct MixedWorkflow;
2297
2298 impl WorkflowHandler for MixedWorkflow {
2299 fn name(&self) -> &str {
2300 "mixed-workflow"
2301 }
2302
2303 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2304 Box::pin(async move {
2305 ctx.shell("build", ShellConfig::new("echo built")).await?;
2306 let op = FakeGitlabOp {
2307 project_id: 456,
2308 title: "Deploy done".to_string(),
2309 };
2310 let result = ctx.operation("notify-gitlab", &op).await?;
2311 assert_eq!(result.output["issue_id"], 42);
2312 Ok(())
2313 })
2314 }
2315 }
2316
2317 #[tokio::test]
2318 async fn operation_step_happy_path() {
2319 let mut engine = create_test_engine();
2320 engine.register(OperationWorkflow).unwrap();
2321
2322 let run = engine
2323 .run_handler("operation-workflow", TriggerKind::Manual, json!({}))
2324 .await
2325 .unwrap()
2326 .run;
2327
2328 assert_eq!(run.status.state, RunStatus::Completed);
2329
2330 let steps = engine.store().list_steps(run.id).await.unwrap();
2331
2332 assert_eq!(steps.len(), 1);
2333 assert_eq!(steps[0].name, "create-issue");
2334 assert_eq!(steps[0].kind, StepKind::Custom("gitlab".to_string()));
2335 assert_eq!(
2336 steps[0].status.state,
2337 ironflow_store::models::StepStatus::Completed
2338 );
2339
2340 let output = steps[0].output.as_ref().unwrap();
2341 assert_eq!(output["issue_id"], 42);
2342 assert_eq!(output["project_id"], 123);
2343
2344 let input = steps[0].input.as_ref().unwrap();
2345 assert_eq!(input["project_id"], 123);
2346 assert_eq!(input["title"], "Bug report");
2347 }
2348
2349 #[tokio::test]
2350 async fn operation_step_failure_marks_run_failed() {
2351 let mut engine = create_test_engine();
2352 engine.register(FailingOperationWorkflow).unwrap();
2353
2354 let result = engine
2355 .run_handler("failing-operation-workflow", TriggerKind::Manual, json!({}))
2356 .await;
2357
2358 assert!(result.is_err());
2359 }
2360
2361 #[tokio::test]
2362 async fn operation_mixed_with_shell_steps() {
2363 let mut engine = create_test_engine();
2364 engine.register(MixedWorkflow).unwrap();
2365
2366 let run = engine
2367 .run_handler("mixed-workflow", TriggerKind::Manual, json!({}))
2368 .await
2369 .unwrap()
2370 .run;
2371
2372 assert_eq!(run.status.state, RunStatus::Completed);
2373
2374 let steps = engine.store().list_steps(run.id).await.unwrap();
2375
2376 assert_eq!(steps.len(), 2);
2377 assert_eq!(steps[0].kind, StepKind::Shell);
2378 assert_eq!(steps[1].kind, StepKind::Custom("gitlab".to_string()));
2379 assert_eq!(steps[0].position, 0);
2380 assert_eq!(steps[1].position, 1);
2381 }
2382
2383 use crate::config::ApprovalConfig;
2388
2389 struct SingleApprovalWorkflow;
2390
2391 impl WorkflowHandler for SingleApprovalWorkflow {
2392 fn name(&self) -> &str {
2393 "single-approval"
2394 }
2395
2396 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2397 Box::pin(async move {
2398 ctx.shell("build", ShellConfig::new("echo built")).await?;
2399 ctx.approval("gate", ApprovalConfig::new("OK?")).await?;
2400 ctx.shell("deploy", ShellConfig::new("echo deployed"))
2401 .await?;
2402 Ok(())
2403 })
2404 }
2405 }
2406
2407 struct DoubleApprovalWorkflow;
2408
2409 impl WorkflowHandler for DoubleApprovalWorkflow {
2410 fn name(&self) -> &str {
2411 "double-approval"
2412 }
2413
2414 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2415 Box::pin(async move {
2416 ctx.shell("build", ShellConfig::new("echo built")).await?;
2417 ctx.approval("staging-gate", ApprovalConfig::new("Deploy staging?"))
2418 .await?;
2419 ctx.shell("deploy-staging", ShellConfig::new("echo staging"))
2420 .await?;
2421 ctx.approval("prod-gate", ApprovalConfig::new("Deploy prod?"))
2422 .await?;
2423 ctx.shell("deploy-prod", ShellConfig::new("echo prod"))
2424 .await?;
2425 Ok(())
2426 })
2427 }
2428 }
2429
2430 #[tokio::test]
2431 async fn approval_pauses_run() {
2432 let mut engine = create_test_engine();
2433 engine.register(SingleApprovalWorkflow).unwrap();
2434
2435 let run = engine
2436 .run_handler("single-approval", TriggerKind::Manual, json!({}))
2437 .await
2438 .unwrap()
2439 .run;
2440
2441 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
2442
2443 let steps = engine.store().list_steps(run.id).await.unwrap();
2444 assert_eq!(steps.len(), 2); assert_eq!(steps[0].kind, StepKind::Shell);
2446 assert_eq!(steps[0].status.state, StepStatus::Completed);
2447 assert_eq!(steps[1].kind, StepKind::Approval);
2448 assert_eq!(steps[1].status.state, StepStatus::AwaitingApproval);
2449 }
2450
2451 #[tokio::test]
2452 async fn approval_resume_completes_run() {
2453 let mut engine = create_test_engine();
2454 engine.register(SingleApprovalWorkflow).unwrap();
2455
2456 let run = engine
2458 .run_handler("single-approval", TriggerKind::Manual, json!({}))
2459 .await
2460 .unwrap()
2461 .run;
2462 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
2463
2464 engine
2466 .store()
2467 .update_run_status(run.id, RunStatus::Running)
2468 .await
2469 .unwrap();
2470
2471 let resumed = engine.resume_run(run.id).await.unwrap().run;
2473 assert_eq!(resumed.status.state, RunStatus::Completed);
2474
2475 let steps = engine.store().list_steps(run.id).await.unwrap();
2476 assert_eq!(steps.len(), 3); assert_eq!(steps[0].name, "build");
2478 assert_eq!(steps[0].status.state, StepStatus::Completed);
2479 assert_eq!(steps[1].name, "gate");
2480 assert_eq!(steps[1].kind, StepKind::Approval);
2481 assert_eq!(steps[1].status.state, StepStatus::Completed);
2482 assert_eq!(steps[2].name, "deploy");
2483 assert_eq!(steps[2].status.state, StepStatus::Completed);
2484 }
2485
2486 #[tokio::test]
2487 async fn double_approval_two_resumes() {
2488 let mut engine = create_test_engine();
2489 engine.register(DoubleApprovalWorkflow).unwrap();
2490
2491 let run = engine
2493 .run_handler("double-approval", TriggerKind::Manual, json!({}))
2494 .await
2495 .unwrap()
2496 .run;
2497 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
2498
2499 let steps = engine.store().list_steps(run.id).await.unwrap();
2500 assert_eq!(steps.len(), 2); engine
2504 .store()
2505 .update_run_status(run.id, RunStatus::Running)
2506 .await
2507 .unwrap();
2508
2509 let resumed = engine.resume_run(run.id).await.unwrap().run;
2510 assert_eq!(resumed.status.state, RunStatus::AwaitingApproval);
2511
2512 let steps = engine.store().list_steps(run.id).await.unwrap();
2513 assert_eq!(steps.len(), 4); engine
2517 .store()
2518 .update_run_status(run.id, RunStatus::Running)
2519 .await
2520 .unwrap();
2521
2522 let final_run = engine.resume_run(run.id).await.unwrap().run;
2523 assert_eq!(final_run.status.state, RunStatus::Completed);
2524
2525 let steps = engine.store().list_steps(run.id).await.unwrap();
2526 assert_eq!(steps.len(), 5);
2527 assert_eq!(steps[0].name, "build");
2528 assert_eq!(steps[1].name, "staging-gate");
2529 assert_eq!(steps[2].name, "deploy-staging");
2530 assert_eq!(steps[3].name, "prod-gate");
2531 assert_eq!(steps[4].name, "deploy-prod");
2532
2533 for step in &steps {
2534 assert_eq!(step.status.state, StepStatus::Completed);
2535 }
2536 }
2537
2538 use ironflow_store::models::{NewStep, StepUpdate, step_trace_id};
2543
2544 async fn create_step_with_status(
2545 store: &Arc<dyn Store>,
2546 run_id: Uuid,
2547 name: &str,
2548 position: u32,
2549 status: StepStatus,
2550 ) -> ironflow_store::models::Step {
2551 let step = store
2552 .create_step(NewStep {
2553 run_id,
2554 trace_id: step_trace_id(run_id, name, position),
2555 name: name.to_string(),
2556 kind: StepKind::Shell,
2557 position,
2558 input: None,
2559 is_error_handler: false,
2560 })
2561 .await
2562 .unwrap();
2563
2564 match status {
2565 StepStatus::Pending => {}
2566 StepStatus::Running => {
2567 store
2568 .update_step(
2569 step.id,
2570 StepUpdate {
2571 status: Some(StepStatus::Running),
2572 ..StepUpdate::default()
2573 },
2574 )
2575 .await
2576 .unwrap();
2577 }
2578 StepStatus::Completed => {
2579 store
2580 .update_step(
2581 step.id,
2582 StepUpdate {
2583 status: Some(StepStatus::Running),
2584 ..StepUpdate::default()
2585 },
2586 )
2587 .await
2588 .unwrap();
2589 store
2590 .update_step(
2591 step.id,
2592 StepUpdate {
2593 status: Some(StepStatus::Completed),
2594 ..StepUpdate::default()
2595 },
2596 )
2597 .await
2598 .unwrap();
2599 }
2600 StepStatus::AwaitingApproval => {
2601 store
2602 .update_step(
2603 step.id,
2604 StepUpdate {
2605 status: Some(StepStatus::Running),
2606 ..StepUpdate::default()
2607 },
2608 )
2609 .await
2610 .unwrap();
2611 store
2612 .update_step(
2613 step.id,
2614 StepUpdate {
2615 status: Some(StepStatus::AwaitingApproval),
2616 ..StepUpdate::default()
2617 },
2618 )
2619 .await
2620 .unwrap();
2621 }
2622 _ => panic!("unsupported status for test helper: {status}"),
2623 }
2624
2625 store.get_step(step.id).await.unwrap().unwrap()
2626 }
2627
2628 #[tokio::test]
2629 async fn fail_orphaned_steps_marks_running_as_failed() {
2630 let engine = create_test_engine();
2631 let run = engine
2632 .store()
2633 .create_run(NewRun {
2634 created_by: None,
2635 workflow_name: "test".to_string(),
2636 trigger: TriggerKind::Manual,
2637 payload: json!({}),
2638 max_retries: 0,
2639 handler_version: None,
2640 labels: HashMap::new(),
2641 scheduled_at: None,
2642 idempotency_key: None,
2643 max_cost_usd: None,
2644 })
2645 .await
2646 .unwrap()
2647 .into_run();
2648
2649 let step = create_step_with_status(
2650 engine.store(),
2651 run.id,
2652 "running-step",
2653 0,
2654 StepStatus::Running,
2655 )
2656 .await;
2657
2658 engine
2659 .fail_orphaned_steps(run.id, "parent run timed out")
2660 .await
2661 .unwrap();
2662
2663 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
2664 assert_eq!(updated.status.state, StepStatus::Failed);
2665 assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
2666 assert!(updated.completed_at.is_some());
2667 }
2668
2669 #[tokio::test]
2670 async fn fail_orphaned_steps_marks_pending_as_skipped() {
2671 let engine = create_test_engine();
2672 let run = engine
2673 .store()
2674 .create_run(NewRun {
2675 created_by: None,
2676 workflow_name: "test".to_string(),
2677 trigger: TriggerKind::Manual,
2678 payload: json!({}),
2679 max_retries: 0,
2680 handler_version: None,
2681 labels: HashMap::new(),
2682 scheduled_at: None,
2683 idempotency_key: None,
2684 max_cost_usd: None,
2685 })
2686 .await
2687 .unwrap()
2688 .into_run();
2689
2690 let step = create_step_with_status(
2691 engine.store(),
2692 run.id,
2693 "pending-step",
2694 0,
2695 StepStatus::Pending,
2696 )
2697 .await;
2698
2699 engine
2700 .fail_orphaned_steps(run.id, "parent run timed out")
2701 .await
2702 .unwrap();
2703
2704 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
2705 assert_eq!(updated.status.state, StepStatus::Skipped);
2706 assert!(updated.error.is_none());
2707 assert!(updated.completed_at.is_some());
2708 }
2709
2710 #[tokio::test]
2711 async fn fail_orphaned_steps_marks_awaiting_approval_as_failed() {
2712 let engine = create_test_engine();
2713 let run = engine
2714 .store()
2715 .create_run(NewRun {
2716 created_by: None,
2717 workflow_name: "test".to_string(),
2718 trigger: TriggerKind::Manual,
2719 payload: json!({}),
2720 max_retries: 0,
2721 handler_version: None,
2722 labels: HashMap::new(),
2723 scheduled_at: None,
2724 idempotency_key: None,
2725 max_cost_usd: None,
2726 })
2727 .await
2728 .unwrap()
2729 .into_run();
2730
2731 let step = create_step_with_status(
2732 engine.store(),
2733 run.id,
2734 "approval-step",
2735 0,
2736 StepStatus::AwaitingApproval,
2737 )
2738 .await;
2739
2740 engine
2741 .fail_orphaned_steps(run.id, "parent run timed out")
2742 .await
2743 .unwrap();
2744
2745 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
2746 assert_eq!(updated.status.state, StepStatus::Failed);
2747 assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
2748 assert!(updated.completed_at.is_some());
2749 }
2750
2751 #[tokio::test]
2752 async fn fail_orphaned_steps_skips_terminal_steps() {
2753 let engine = create_test_engine();
2754 let run = engine
2755 .store()
2756 .create_run(NewRun {
2757 created_by: None,
2758 workflow_name: "test".to_string(),
2759 trigger: TriggerKind::Manual,
2760 payload: json!({}),
2761 max_retries: 0,
2762 handler_version: None,
2763 labels: HashMap::new(),
2764 scheduled_at: None,
2765 idempotency_key: None,
2766 max_cost_usd: None,
2767 })
2768 .await
2769 .unwrap()
2770 .into_run();
2771
2772 let completed_step =
2773 create_step_with_status(engine.store(), run.id, "done", 0, StepStatus::Completed).await;
2774 let running_step =
2775 create_step_with_status(engine.store(), run.id, "in-flight", 1, StepStatus::Running)
2776 .await;
2777
2778 engine
2779 .fail_orphaned_steps(run.id, "parent run timed out")
2780 .await
2781 .unwrap();
2782
2783 let completed = engine
2784 .store()
2785 .get_step(completed_step.id)
2786 .await
2787 .unwrap()
2788 .unwrap();
2789 assert_eq!(completed.status.state, StepStatus::Completed);
2790
2791 let failed = engine
2792 .store()
2793 .get_step(running_step.id)
2794 .await
2795 .unwrap()
2796 .unwrap();
2797 assert_eq!(failed.status.state, StepStatus::Failed);
2798 }
2799
2800 #[tokio::test]
2801 async fn fail_orphaned_steps_mixed_states() {
2802 let engine = create_test_engine();
2803 let run = engine
2804 .store()
2805 .create_run(NewRun {
2806 created_by: None,
2807 workflow_name: "test".to_string(),
2808 trigger: TriggerKind::Manual,
2809 payload: json!({}),
2810 max_retries: 0,
2811 handler_version: None,
2812 labels: HashMap::new(),
2813 scheduled_at: None,
2814 idempotency_key: None,
2815 max_cost_usd: None,
2816 })
2817 .await
2818 .unwrap()
2819 .into_run();
2820
2821 let s_completed =
2822 create_step_with_status(engine.store(), run.id, "step-1", 0, StepStatus::Completed)
2823 .await;
2824 let s_running =
2825 create_step_with_status(engine.store(), run.id, "step-2", 1, StepStatus::Running).await;
2826 let s_pending =
2827 create_step_with_status(engine.store(), run.id, "step-3", 2, StepStatus::Pending).await;
2828
2829 engine.fail_orphaned_steps(run.id, "timeout").await.unwrap();
2830
2831 let r_completed = engine
2832 .store()
2833 .get_step(s_completed.id)
2834 .await
2835 .unwrap()
2836 .unwrap();
2837 assert_eq!(r_completed.status.state, StepStatus::Completed);
2838
2839 let r_running = engine
2840 .store()
2841 .get_step(s_running.id)
2842 .await
2843 .unwrap()
2844 .unwrap();
2845 assert_eq!(r_running.status.state, StepStatus::Failed);
2846 assert_eq!(r_running.error.as_deref(), Some("timeout"));
2847
2848 let r_pending = engine
2849 .store()
2850 .get_step(s_pending.id)
2851 .await
2852 .unwrap()
2853 .unwrap();
2854 assert_eq!(r_pending.status.state, StepStatus::Skipped);
2855 assert!(r_pending.error.is_none());
2856 }
2857
2858 #[tokio::test]
2859 async fn fail_orphaned_steps_no_steps_is_noop() {
2860 let engine = create_test_engine();
2861 let run = engine
2862 .store()
2863 .create_run(NewRun {
2864 created_by: None,
2865 workflow_name: "test".to_string(),
2866 trigger: TriggerKind::Manual,
2867 payload: json!({}),
2868 max_retries: 0,
2869 handler_version: None,
2870 labels: HashMap::new(),
2871 scheduled_at: None,
2872 idempotency_key: None,
2873 max_cost_usd: None,
2874 })
2875 .await
2876 .unwrap()
2877 .into_run();
2878
2879 let result = engine.fail_orphaned_steps(run.id, "timeout").await;
2880 assert!(result.is_ok());
2881 }
2882
2883 #[tokio::test]
2884 async fn fail_orphaned_steps_preserves_existing_error() {
2885 let engine = create_test_engine();
2886 let run = engine
2887 .store()
2888 .create_run(NewRun {
2889 created_by: None,
2890 workflow_name: "test".to_string(),
2891 trigger: TriggerKind::Manual,
2892 payload: json!({}),
2893 max_retries: 0,
2894 handler_version: None,
2895 labels: HashMap::new(),
2896 scheduled_at: None,
2897 idempotency_key: None,
2898 max_cost_usd: None,
2899 })
2900 .await
2901 .unwrap()
2902 .into_run();
2903
2904 let step_with_error = create_step_with_status(
2905 engine.store(),
2906 run.id,
2907 "already-errored",
2908 0,
2909 StepStatus::Running,
2910 )
2911 .await;
2912
2913 engine
2914 .store()
2915 .update_step(
2916 step_with_error.id,
2917 StepUpdate {
2918 error: Some("real error from provider".to_string()),
2919 ..StepUpdate::default()
2920 },
2921 )
2922 .await
2923 .unwrap();
2924
2925 let step_no_error = create_step_with_status(
2926 engine.store(),
2927 run.id,
2928 "no-error-yet",
2929 1,
2930 StepStatus::Running,
2931 )
2932 .await;
2933
2934 engine
2935 .fail_orphaned_steps(run.id, "parent run failed")
2936 .await
2937 .unwrap();
2938
2939 let updated_with = engine
2940 .store()
2941 .get_step(step_with_error.id)
2942 .await
2943 .unwrap()
2944 .unwrap();
2945 assert_eq!(updated_with.status.state, StepStatus::Failed);
2946 assert_eq!(
2947 updated_with.error.as_deref(),
2948 Some("real error from provider"),
2949 );
2950
2951 let updated_without = engine
2952 .store()
2953 .get_step(step_no_error.id)
2954 .await
2955 .unwrap()
2956 .unwrap();
2957 assert_eq!(updated_without.status.state, StepStatus::Failed);
2958 assert_eq!(updated_without.error.as_deref(), Some("parent run failed"),);
2959 }
2960}