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