1use std::collections::{HashMap, HashSet};
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, to_value};
17use tokio::spawn;
18use tracing::{error, info, warn};
19use uuid::Uuid;
20
21use ironflow_core::error::OperationError;
22#[cfg(feature = "prometheus")]
23use ironflow_core::metric_names::{
24 RUN_BUDGET_EXCEEDED_TOTAL, RUN_COST_USD, RUN_DURATION_SECONDS, RUNS_ACTIVE, RUNS_TOTAL,
25};
26use ironflow_core::provider::{AgentProvider, LABEL_ROOT_RUN_ID};
27use ironflow_store::error::StoreError;
28use ironflow_store::models::{
29 ConcurrencyLimit, NewRun, NewSignal, Run, RunActor, RunCreation, RunFilter, RunStatus,
30 RunUpdate, SignalInsert, SignalStepResolution, StepStatus, StepUpdate, TriggerKind,
31 validate_concurrency_limits,
32};
33use ironflow_store::store::Store;
34#[cfg(feature = "prometheus")]
35use metrics::{counter, gauge, histogram};
36
37use crate::artifact::ArtifactSink;
38use crate::budget::{BudgetConfig, month_start};
39use crate::context::{PARENT_RUN_ID_LABEL, WorkflowContext, interrupt_running_steps};
40use crate::error::EngineError;
41use crate::executor::{StepInterceptor, StepResult};
42use crate::guard::{WorkflowGuardConfig, new_shared_guard_state};
43use crate::handler::{WorkflowHandler, WorkflowInfo};
44use crate::log_sender::LogSender;
45use crate::notify::{
46 ApprovalRequestedEvent, Event, EventPublisher, EventSubscriber, RunBudgetExceededEvent,
47 RunFailedEvent, RunStatusChangedEvent, SignalAwaitedEvent, SignalReceivedEvent,
48 WorkflowEventBus,
49};
50use crate::plan::{
51 ExecutionPlan, PlanOptions, PlanRecorder, SharedPlanRecorder, estimate_durations, lock_plan,
52};
53use crate::retry_policy::{backoff_for_retry, is_run_retryable};
54use crate::schedule::CronSchedule;
55use crate::signal::{
56 Signal, SignalDelivery, SignalRejected, SignalResumed, received_output, validate_step_payload,
57};
58use ironflow_core::decision::DecisionProvider;
59
60#[derive(Debug, Clone)]
79pub struct WorkflowResult {
80 pub run: Run,
82 pub steps: Vec<StepResult>,
84}
85
86#[derive(Debug, Clone, Default)]
105pub struct EnqueueOptions {
106 pub max_retries: u32,
108 pub labels: HashMap<String, String>,
110 pub scheduled_at: Option<DateTime<Utc>>,
113 pub max_cost_usd: Option<Decimal>,
117 pub created_by: Option<RunActor>,
120 pub idempotency_key: Option<String>,
126 pub concurrency_key: Option<String>,
132 pub concurrency_limits: Vec<ConcurrencyLimit>,
140}
141
142#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
153pub enum ExecutionMode {
154 #[default]
158 Local,
159 Workers,
162}
163
164pub struct Engine {
205 store: Arc<dyn Store>,
206 provider: Arc<dyn AgentProvider>,
207 handlers: HashMap<String, Arc<dyn WorkflowHandler>>,
208 event_publisher: EventPublisher,
209 log_sender: Option<LogSender>,
210 budget: BudgetConfig,
211 artifact_sink: Option<Arc<dyn ArtifactSink>>,
212 guard_config: Option<WorkflowGuardConfig>,
213 event_bus: Option<WorkflowEventBus>,
214 decision_provider: Option<Arc<dyn DecisionProvider>>,
215 step_interceptor: Option<Arc<dyn StepInterceptor>>,
216 execution_mode: ExecutionMode,
217}
218
219fn validate_category(handler_name: &str, category: &str) -> Result<(), EngineError> {
229 let reject = |reason: &str| {
230 Err(EngineError::InvalidWorkflow(format!(
231 "handler '{handler_name}' has invalid category '{category}': {reason}"
232 )))
233 };
234
235 if category.is_empty() {
236 return reject("empty category");
237 }
238 if category.starts_with('/') {
239 return reject("leading '/'");
240 }
241 if category.ends_with('/') {
242 return reject("trailing '/'");
243 }
244 for segment in category.split('/') {
245 if segment.is_empty() {
246 return reject("empty segment (double '/')");
247 }
248 if segment.trim().is_empty() {
249 return reject("whitespace-only segment");
250 }
251 }
252 Ok(())
253}
254
255fn chain_label(run: &Run, key: &str) -> Option<Uuid> {
260 if !matches!(run.trigger, TriggerKind::Workflow) {
261 return None;
262 }
263 let id = Uuid::parse_str(run.labels.get(key)?).ok()?;
264 (id != run.id).then_some(id)
265}
266
267pub(crate) fn chain_root(run: &Run) -> Option<Uuid> {
272 chain_label(run, LABEL_ROOT_RUN_ID)
273}
274
275fn chain_parent(run: &Run) -> Option<Uuid> {
277 chain_label(run, PARENT_RUN_ID_LABEL)
278}
279
280impl Engine {
281 pub fn new(store: Arc<dyn Store>, provider: Arc<dyn AgentProvider>) -> Self {
297 Self {
298 store,
299 provider,
300 handlers: HashMap::new(),
301 event_publisher: EventPublisher::new(),
302 log_sender: None,
303 budget: BudgetConfig::new(),
304 artifact_sink: None,
305 guard_config: None,
306 event_bus: None,
307 decision_provider: None,
308 step_interceptor: None,
309 execution_mode: ExecutionMode::default(),
310 }
311 }
312
313 pub fn with_decision_provider(mut self, provider: Arc<dyn DecisionProvider>) -> Self {
334 self.decision_provider = Some(provider);
335 self
336 }
337
338 pub fn with_step_interceptor(mut self, interceptor: Arc<dyn StepInterceptor>) -> Self {
362 self.step_interceptor = Some(interceptor);
363 self
364 }
365
366 pub fn step_interceptor(&self) -> Option<&Arc<dyn StepInterceptor>> {
368 self.step_interceptor.as_ref()
369 }
370
371 pub fn with_budget_config(mut self, budget: BudgetConfig) -> Self {
392 self.budget = budget;
393 self
394 }
395
396 pub fn budget_config(&self) -> &BudgetConfig {
398 &self.budget
399 }
400
401 pub fn with_guard_config(mut self, config: WorkflowGuardConfig) -> Self {
423 self.guard_config = Some(config);
424 self
425 }
426
427 pub fn guard_config(&self) -> Option<&WorkflowGuardConfig> {
429 self.guard_config.as_ref()
430 }
431
432 pub fn with_execution_mode(mut self, mode: ExecutionMode) -> Self {
453 self.execution_mode = mode;
454 self
455 }
456
457 pub fn execution_mode(&self) -> ExecutionMode {
459 self.execution_mode
460 }
461
462 pub fn set_log_sender(&mut self, sender: LogSender) {
468 self.log_sender = Some(sender);
469 }
470
471 pub fn set_artifact_sink(&mut self, sink: Arc<dyn ArtifactSink>) {
489 self.artifact_sink = Some(sink);
490 }
491
492 pub fn artifact_sink(&self) -> Option<&Arc<dyn ArtifactSink>> {
494 self.artifact_sink.as_ref()
495 }
496
497 pub fn set_event_bus(&mut self, bus: WorkflowEventBus) {
514 self.event_bus = Some(bus);
515 }
516
517 pub fn event_bus(&self) -> Option<&WorkflowEventBus> {
519 self.event_bus.as_ref()
520 }
521
522 pub fn store(&self) -> &Arc<dyn Store> {
524 &self.store
525 }
526
527 pub fn provider(&self) -> &Arc<dyn AgentProvider> {
529 &self.provider
530 }
531
532 fn build_context(&self, run: &Run) -> WorkflowContext {
541 let handlers = self.handlers.clone();
542 let resolver: crate::context::HandlerResolver =
543 Arc::new(move |name: &str| handlers.get(name).cloned());
544 let mut ctx = WorkflowContext::with_handler_resolver(
545 run.id,
546 run.workflow_name.clone(),
547 self.store.clone(),
548 self.provider.clone(),
549 resolver,
550 );
551 ctx.carry_over_run_totals(run.retry_count + 1, run.cost_usd, run.duration_ms);
552 ctx.set_max_cost_usd(run.max_cost_usd);
553 ctx.set_run_created_at(run.created_at);
554 if let Some(ref sender) = self.log_sender {
555 ctx.set_log_sender(sender.clone());
556 }
557 if let Some(ref sink) = self.artifact_sink {
558 ctx.set_artifact_sink(sink.clone());
559 }
560 if let Some(ref bus) = self.event_bus {
561 ctx.set_event_bus(bus.clone());
562 }
563 if let Some(ref provider) = self.decision_provider {
564 ctx.set_decision_provider(provider.clone());
565 }
566 if let Some(ref interceptor) = self.step_interceptor {
567 ctx.set_step_interceptor(interceptor.clone());
568 }
569 ctx
570 }
571
572 fn build_context_with_guard(
578 &self,
579 run: &Run,
580 handler: &dyn WorkflowHandler,
581 ) -> WorkflowContext {
582 let mut ctx = self.build_context(run);
583 let guard_config = handler.guard_config().or_else(|| self.guard_config.clone());
584 if let Some(config) = guard_config {
585 ctx.set_guard(config, new_shared_guard_state());
586 }
587 ctx
588 }
589
590 async fn check_monthly_quota(&self, workflow_name: &str) -> Result<(), EngineError> {
601 let Some(limit) = self.budget.monthly_cost_limit_usd else {
602 return Ok(());
603 };
604
605 let stats = self
606 .store
607 .get_stats(RunFilter {
608 created_after: Some(month_start(Utc::now())),
609 ..RunFilter::default()
610 })
611 .await?;
612
613 if stats.total_cost_usd < limit {
614 return Ok(());
615 }
616
617 warn!(
618 workflow = %workflow_name,
619 limit_usd = %limit,
620 spent_usd = %stats.total_cost_usd,
621 "monthly cost quota exhausted, refusing new run"
622 );
623
624 #[cfg(feature = "prometheus")]
625 counter!(
626 RUN_BUDGET_EXCEEDED_TOTAL,
627 "workflow" => workflow_name.to_string(),
628 "scope" => "monthly",
629 )
630 .increment(1);
631
632 Err(EngineError::MonthlyBudgetExceeded {
633 limit_usd: limit,
634 spent_usd: stats.total_cost_usd,
635 })
636 }
637
638 pub fn register(&mut self, handler: impl WorkflowHandler + 'static) -> Result<(), EngineError> {
682 let name = handler.name().to_string();
683 if self.handlers.contains_key(&name) {
684 return Err(EngineError::InvalidWorkflow(format!(
685 "handler '{}' already registered",
686 name
687 )));
688 }
689 if let Some(category) = handler.category() {
690 validate_category(&name, category)?;
691 }
692 self.handlers.insert(name, Arc::new(handler));
693 Ok(())
694 }
695
696 pub fn register_boxed(&mut self, handler: Box<dyn WorkflowHandler>) -> Result<(), EngineError> {
703 let name = handler.name().to_string();
704 if self.handlers.contains_key(&name) {
705 return Err(EngineError::InvalidWorkflow(format!(
706 "handler '{}' already registered",
707 name
708 )));
709 }
710 if let Some(category) = handler.category() {
711 validate_category(&name, category)?;
712 }
713 self.handlers.insert(name, Arc::from(handler));
714 Ok(())
715 }
716
717 pub fn get_handler(&self, name: &str) -> Option<&Arc<dyn WorkflowHandler>> {
719 self.handlers.get(name)
720 }
721
722 pub fn handler_names(&self) -> Vec<&str> {
724 self.handlers.keys().map(|s| s.as_str()).collect()
725 }
726
727 pub fn handler_info(&self, name: &str) -> Option<WorkflowInfo> {
729 self.handlers.get(name).map(|h| h.describe())
730 }
731
732 pub fn scheduled_handlers(&self) -> Vec<(&str, &CronSchedule)> {
756 self.handlers
757 .iter()
758 .filter_map(|(name, handler)| handler.schedule().map(|sched| (name.as_str(), sched)))
759 .collect()
760 }
761
762 pub fn subscribe(
787 &mut self,
788 subscriber: impl EventSubscriber + 'static,
789 event_types: &[&'static str],
790 ) {
791 self.event_publisher.subscribe(subscriber, event_types);
792 }
793
794 pub fn event_publisher(&self) -> &EventPublisher {
799 &self.event_publisher
800 }
801
802 #[tracing::instrument(name = "engine.run_handler", skip_all, fields(workflow = %handler_name))]
832 pub async fn run_handler(
833 &self,
834 handler_name: &str,
835 trigger: TriggerKind,
836 payload: Value,
837 ) -> Result<WorkflowResult, EngineError> {
838 let handler = self
839 .handlers
840 .get(handler_name)
841 .ok_or_else(|| {
842 EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
843 })?
844 .clone();
845
846 self.check_monthly_quota(handler_name).await?;
847
848 let handler_version = handler.version().map(str::to_string);
849 let max_cost_usd = self
850 .budget
851 .resolve_run_cap(None, handler.default_max_cost_usd());
852 let run = self
853 .store
854 .create_run(NewRun {
855 created_by: None,
856 workflow_name: handler_name.to_string(),
857 trigger,
858 payload,
859 max_retries: 0,
860 handler_version,
861 labels: handler.default_labels(),
862 scheduled_at: None,
863 idempotency_key: None,
864 concurrency_key: None,
865 concurrency_limits: Vec::new(),
866 max_cost_usd,
867 })
868 .await?
869 .into_run();
870
871 let run_id = run.id;
872 info!(run_id = %run_id, handler_version = run.handler_version.as_deref().unwrap_or(""), "run created");
873
874 self.store
875 .update_run_status(run_id, RunStatus::Running)
876 .await?;
877
878 #[cfg(feature = "prometheus")]
879 gauge!(RUNS_ACTIVE, "workflow" => handler_name.to_string()).increment(1.0);
880
881 let run_start = Instant::now();
882 let mut ctx = self.build_context_with_guard(&run, handler.as_ref());
883
884 let result = handler.execute(&mut ctx).await;
885 self.finalize_run(run_id, handler_name, result, &ctx, run_start, run.labels)
886 .await
887 }
888
889 #[tracing::instrument(name = "engine.plan_handler", skip_all, fields(workflow = %handler_name))]
931 pub async fn plan_handler(
932 &self,
933 handler_name: &str,
934 payload: Value,
935 options: PlanOptions,
936 ) -> Result<ExecutionPlan, EngineError> {
937 if options.max_depth == 0 {
938 return Err(EngineError::InvalidWorkflow(
939 "max_depth must be at least 1".to_string(),
940 ));
941 }
942
943 let handler = self
944 .handlers
945 .get(handler_name)
946 .ok_or_else(|| {
947 EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
948 })?
949 .clone();
950
951 let estimates = if options.estimate_durations {
952 estimate_durations(&self.store, handler_name, options.sample_runs).await?
953 } else {
954 HashMap::new()
955 };
956
957 let shared: SharedPlanRecorder = Arc::new(Mutex::new(PlanRecorder::new(
958 handler_name.to_string(),
959 payload,
960 options.max_depth,
961 estimates,
962 )));
963
964 let handlers = self.handlers.clone();
967 let resolver: crate::context::HandlerResolver =
968 Arc::new(move |name: &str| handlers.get(name).cloned());
969 let mut ctx = WorkflowContext::with_handler_resolver(
970 Uuid::now_v7(),
971 handler_name.to_string(),
972 self.store.clone(),
973 self.provider.clone(),
974 resolver,
975 );
976 ctx.set_plan(shared.clone());
977
978 if let Err(err) = handler.execute(&mut ctx).await {
979 lock_plan(&shared).fail(err.to_string());
980 }
981 drop(ctx);
982
983 let plan = match Arc::try_unwrap(shared) {
984 Ok(mutex) => mutex
985 .into_inner()
986 .unwrap_or_else(|poisoned| poisoned.into_inner())
987 .into_plan(),
988 Err(shared) => lock_plan(&shared).snapshot(),
989 };
990
991 info!(
992 workflow = %handler_name,
993 steps = plan.steps.len(),
994 truncated = plan.truncated,
995 "execution plan built"
996 );
997
998 Ok(plan)
999 }
1000
1001 #[tracing::instrument(name = "engine.enqueue_handler", skip_all, fields(workflow = %handler_name))]
1012 pub async fn enqueue_handler(
1013 &self,
1014 handler_name: &str,
1015 trigger: TriggerKind,
1016 payload: Value,
1017 max_retries: u32,
1018 ) -> Result<Run, EngineError> {
1019 self.enqueue_handler_with_options(
1020 handler_name,
1021 trigger,
1022 payload,
1023 EnqueueOptions {
1024 max_retries,
1025 ..Default::default()
1026 },
1027 )
1028 .await
1029 .map(RunCreation::into_run)
1030 }
1031
1032 #[tracing::instrument(name = "engine.enqueue_handler_with_options", skip_all, fields(workflow = %handler_name))]
1080 pub async fn enqueue_handler_with_options(
1081 &self,
1082 handler_name: &str,
1083 trigger: TriggerKind,
1084 payload: Value,
1085 options: EnqueueOptions,
1086 ) -> Result<RunCreation, EngineError> {
1087 let EnqueueOptions {
1088 max_retries,
1089 labels,
1090 scheduled_at,
1091 max_cost_usd,
1092 created_by,
1093 idempotency_key,
1094 concurrency_key,
1095 concurrency_limits,
1096 } = options;
1097
1098 validate_concurrency_limits(&concurrency_limits)
1101 .map_err(EngineError::InvalidConcurrencyLimit)?;
1102
1103 let handler = self.handlers.get(handler_name).ok_or_else(|| {
1104 EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
1105 })?;
1106
1107 self.check_monthly_quota(handler_name).await?;
1108
1109 let handler_version = handler.version().map(str::to_string);
1110 let mut merged_labels = handler.default_labels();
1111 merged_labels.extend(labels);
1112 let resolved_cap = self
1113 .budget
1114 .resolve_run_cap(max_cost_usd, handler.default_max_cost_usd());
1115
1116 let creation = self
1117 .store
1118 .create_run(NewRun {
1119 workflow_name: handler_name.to_string(),
1120 trigger,
1121 payload,
1122 max_retries,
1123 handler_version,
1124 labels: merged_labels,
1125 scheduled_at,
1126 created_by,
1127 idempotency_key,
1128 concurrency_key,
1129 concurrency_limits,
1130 max_cost_usd: resolved_cap,
1131 })
1132 .await?;
1133
1134 match &creation {
1135 RunCreation::Created(run) => info!(
1136 run_id = %run.id,
1137 workflow = %handler_name,
1138 max_cost_usd = ?resolved_cap,
1139 "handler run enqueued"
1140 ),
1141 RunCreation::Existing(run) => info!(
1142 run_id = %run.id,
1143 workflow = %handler_name,
1144 "idempotent replay, nothing enqueued"
1145 ),
1146 }
1147
1148 Ok(creation)
1149 }
1150
1151 #[tracing::instrument(name = "engine.execute_handler_run", skip_all, fields(run_id = %run_id))]
1171 pub async fn execute_handler_run(&self, run_id: Uuid) -> Result<WorkflowResult, EngineError> {
1172 let run = self
1173 .store
1174 .get_run(run_id)
1175 .await?
1176 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
1177
1178 if let Some(root_run_id) = chain_root(&run) {
1179 return self.resume_chain(run_id, root_run_id).await;
1180 }
1181
1182 let handler = self
1183 .handlers
1184 .get(&run.workflow_name)
1185 .ok_or_else(|| {
1186 EngineError::InvalidWorkflow(format!(
1187 "no handler registered: {}",
1188 run.workflow_name
1189 ))
1190 })?
1191 .clone();
1192
1193 #[cfg(feature = "prometheus")]
1194 gauge!(RUNS_ACTIVE, "workflow" => run.workflow_name.clone()).increment(1.0);
1195
1196 let run_start = Instant::now();
1197 let mut ctx = self.build_context_with_guard(&run, handler.as_ref());
1198
1199 let result = if handler.is_version_compatible(run.handler_version.as_deref()) {
1212 ctx.load_replay_steps().await?;
1213 self.release_then_execute(run_id, handler.as_ref(), &mut ctx)
1214 .await
1215 } else {
1216 Err(EngineError::HandlerVersionMismatch {
1217 run_id,
1218 workflow_name: run.workflow_name.clone(),
1219 run_version: run
1220 .handler_version
1221 .clone()
1222 .unwrap_or_else(|| "unknown".to_string()),
1223 current_version: handler
1224 .version()
1225 .map(str::to_string)
1226 .unwrap_or_else(|| "unknown".to_string()),
1227 })
1228 };
1229
1230 self.finalize_run(
1231 run_id,
1232 &run.workflow_name,
1233 result,
1234 &ctx,
1235 run_start,
1236 run.labels,
1237 )
1238 .await
1239 }
1240
1241 #[tracing::instrument(name = "engine.execute_run", skip_all, fields(run_id = %run_id))]
1249 pub async fn execute_run(&self, run_id: Uuid) -> Result<WorkflowResult, EngineError> {
1250 self.execute_handler_run(run_id).await
1251 }
1252
1253 #[tracing::instrument(name = "engine.resume_run", skip_all, fields(run_id = %run_id))]
1282 pub async fn resume_run(&self, run_id: Uuid) -> Result<WorkflowResult, EngineError> {
1283 let run = self
1284 .store
1285 .get_run(run_id)
1286 .await?
1287 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
1288
1289 if let Some(root_run_id) = chain_root(&run) {
1290 return self.resume_chain(run_id, root_run_id).await;
1291 }
1292
1293 self.resume_loaded_run(run).await
1294 }
1295
1296 async fn resume_chain(
1304 &self,
1305 child_run_id: Uuid,
1306 root_run_id: Uuid,
1307 ) -> Result<WorkflowResult, EngineError> {
1308 let root = self
1309 .store
1310 .get_run(root_run_id)
1311 .await?
1312 .ok_or(EngineError::Store(StoreError::RunNotFound(root_run_id)))?;
1313
1314 match root.status.state {
1315 RunStatus::AwaitingApproval | RunStatus::Pending => {
1316 self.store
1317 .update_run_status(root_run_id, RunStatus::Running)
1318 .await?;
1319 }
1320 RunStatus::Sleeping => {
1321 self.store
1322 .update_run_status(root_run_id, RunStatus::Pending)
1323 .await?;
1324 self.store
1325 .update_run_status(root_run_id, RunStatus::Running)
1326 .await?;
1327 }
1328 other => {
1329 let reason = format!(
1330 "cannot resume child run {child_run_id}: root run {root_run_id} is {other}"
1331 );
1332 if let Err(err) = self
1333 .fail_or_schedule_retry(child_run_id, &reason, false, None, None)
1334 .await
1335 {
1336 error!(
1337 run_id = %child_run_id,
1338 error = %err,
1339 "failed to fail a child run whose root cannot resume"
1340 );
1341 }
1342 return Err(EngineError::InvalidWorkflow(reason));
1343 }
1344 }
1345
1346 info!(
1347 run_id = %child_run_id,
1348 root_run_id = %root_run_id,
1349 "child run resumed through its root run"
1350 );
1351
1352 let root = self
1353 .store
1354 .get_run(root_run_id)
1355 .await?
1356 .ok_or(EngineError::Store(StoreError::RunNotFound(root_run_id)))?;
1357 self.resume_loaded_run(root).await
1358 }
1359
1360 async fn resume_loaded_run(&self, run: Run) -> Result<WorkflowResult, EngineError> {
1362 let run_id = run.id;
1363 let handler = self
1364 .handlers
1365 .get(&run.workflow_name)
1366 .ok_or_else(|| {
1367 EngineError::InvalidWorkflow(format!(
1368 "no handler registered: {}",
1369 run.workflow_name
1370 ))
1371 })?
1372 .clone();
1373
1374 info!(run_id = %run_id, workflow = %run.workflow_name, "resuming run after approval");
1375
1376 let run_start = Instant::now();
1377 let mut ctx = self.build_context_with_guard(&run, handler.as_ref());
1378
1379 let result = if handler.is_version_compatible(run.handler_version.as_deref()) {
1380 ctx.load_replay_steps().await?;
1381 self.release_then_execute(run_id, handler.as_ref(), &mut ctx)
1382 .await
1383 } else {
1384 Err(EngineError::HandlerVersionMismatch {
1385 run_id,
1386 workflow_name: run.workflow_name.clone(),
1387 run_version: run
1388 .handler_version
1389 .clone()
1390 .unwrap_or_else(|| "unknown".to_string()),
1391 current_version: handler
1392 .version()
1393 .map(str::to_string)
1394 .unwrap_or_else(|| "unknown".to_string()),
1395 })
1396 };
1397
1398 self.finalize_run(
1399 run_id,
1400 &run.workflow_name,
1401 result,
1402 &ctx,
1403 run_start,
1404 run.labels,
1405 )
1406 .await
1407 }
1408
1409 pub async fn deliver_signal(
1455 self: &Arc<Self>,
1456 signal: NewSignal,
1457 ) -> Result<SignalDelivery, EngineError> {
1458 if signal.name.trim().is_empty() {
1459 return Err(EngineError::InvalidSignal(
1460 "signal name must not be empty".to_string(),
1461 ));
1462 }
1463 if signal.key.trim().is_empty() {
1464 return Err(EngineError::InvalidSignal(
1465 "signal key must not be empty".to_string(),
1466 ));
1467 }
1468
1469 let stored = match self.store.insert_signal(signal).await? {
1470 SignalInsert::Created(stored) => stored,
1471 SignalInsert::Duplicate(existing) => {
1472 info!(
1473 signal_id = %existing.id,
1474 signal = %existing.name,
1475 key = %existing.key,
1476 "duplicate signal ignored"
1477 );
1478 return Ok(SignalDelivery {
1479 signal_id: existing.id,
1480 duplicate: true,
1481 resumed: Vec::new(),
1482 rejected: Vec::new(),
1483 });
1484 }
1485 };
1486
1487 let waiters = self
1488 .store
1489 .list_signal_waiters(&stored.name, &stored.key)
1490 .await?;
1491 let mut resumed = Vec::new();
1492 let mut rejected = Vec::new();
1493
1494 for step in waiters {
1495 if let Err(error) = validate_step_payload(step.input.as_ref(), &stored.payload) {
1496 rejected.push(SignalRejected {
1497 run_id: step.run_id,
1498 step_id: step.id,
1499 error,
1500 });
1501 continue;
1502 }
1503
1504 match self
1505 .store
1506 .resolve_signal_step(step.id, received_output(&stored))
1507 .await
1508 {
1509 Ok(SignalStepResolution::Resolved {
1510 run_id,
1511 run_resumed,
1512 }) => {
1513 resumed.push(SignalResumed {
1514 run_id,
1515 step_id: step.id,
1516 });
1517 if run_resumed && self.execution_mode == ExecutionMode::Local {
1518 self.spawn_local_resume(run_id);
1519 }
1520 }
1521 Ok(SignalStepResolution::NotWaiting { .. }) => {}
1523 Err(err) => {
1524 error!(
1525 run_id = %step.run_id,
1526 step_id = %step.id,
1527 error = %err,
1528 "failed to resolve a waiting signal step"
1529 );
1530 rejected.push(SignalRejected {
1531 run_id: step.run_id,
1532 step_id: step.id,
1533 error: err.to_string(),
1534 });
1535 }
1536 }
1537 }
1538
1539 info!(
1540 signal_id = %stored.id,
1541 signal = %stored.name,
1542 key = %stored.key,
1543 resumed = resumed.len(),
1544 rejected = rejected.len(),
1545 "signal received"
1546 );
1547 self.event_publisher
1548 .publish(Event::SignalReceived(SignalReceivedEvent {
1549 signal_id: stored.id,
1550 name: stored.name.clone(),
1551 key: stored.key.clone(),
1552 resumed_runs: resumed.iter().map(|r| r.run_id).collect(),
1553 at: stored.received_at,
1554 }));
1555
1556 Ok(SignalDelivery {
1557 signal_id: stored.id,
1558 duplicate: false,
1559 resumed,
1560 rejected,
1561 })
1562 }
1563
1564 pub async fn send_signal<S: Signal>(
1602 self: &Arc<Self>,
1603 signal: &S,
1604 key: &str,
1605 idempotency_id: Option<&str>,
1606 ) -> Result<SignalDelivery, EngineError> {
1607 let payload = to_value(signal)?;
1608 self.deliver_signal(NewSignal {
1609 name: S::NAME.to_string(),
1610 key: key.to_string(),
1611 payload,
1612 idempotency_id: idempotency_id.map(str::to_string),
1613 })
1614 .await
1615 }
1616
1617 pub(crate) fn spawn_local_resume(self: &Arc<Self>, run_id: Uuid) {
1623 let engine = Arc::clone(self);
1624 spawn(async move {
1625 if let Err(err) = engine
1626 .store
1627 .update_run_status(run_id, RunStatus::Running)
1628 .await
1629 {
1630 error!(run_id = %run_id, error = %err, "failed to restart a woken run");
1631 return;
1632 }
1633 if let Err(err) = engine.resume_run(run_id).await {
1634 error!(run_id = %run_id, error = %err, "failed to resume a woken run");
1635 }
1636 });
1637 }
1638
1639 pub async fn fail_or_schedule_retry(
1687 &self,
1688 run_id: Uuid,
1689 error: &str,
1690 retryable: bool,
1691 cost_usd: Option<Decimal>,
1692 duration_ms: Option<u64>,
1693 ) -> Result<RunStatus, EngineError> {
1694 let run = self
1695 .store
1696 .get_run(run_id)
1697 .await?
1698 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
1699
1700 let has_attempts_left = run.retry_count < run.max_retries;
1701 let update = if retryable && has_attempts_left {
1702 let backoff = backoff_for_retry(run.retry_count);
1703 let scheduled_at = Utc::now() + TimeDelta::milliseconds(backoff.as_millis() as i64);
1704
1705 info!(
1706 run_id = %run_id,
1707 workflow = %run.workflow_name,
1708 attempt = run.retry_count + 1,
1709 max_retries = run.max_retries,
1710 backoff_secs = backoff.as_secs(),
1711 scheduled_at = %scheduled_at,
1712 "run failed, scheduling retry"
1713 );
1714
1715 RunUpdate {
1716 status: Some(RunStatus::Retrying),
1717 error: Some(error.to_string()),
1718 increment_retry: true,
1719 cost_usd,
1720 duration_ms,
1721 scheduled_at: Some(scheduled_at),
1722 ..RunUpdate::default()
1723 }
1724 } else {
1725 RunUpdate {
1726 status: Some(RunStatus::Failed),
1727 error: Some(error.to_string()),
1728 cost_usd,
1729 duration_ms,
1730 completed_at: Some(Utc::now()),
1731 ..RunUpdate::default()
1732 }
1733 };
1734
1735 let status = update.status.unwrap_or(RunStatus::Failed);
1736 self.store.update_run(run_id, update).await?;
1737 self.fail_orphaned_steps(run_id, error).await?;
1738 self.cancel_descendants_of_stopped_run(run_id, error).await;
1741
1742 Ok(status)
1743 }
1744
1745 pub async fn interrupt_running_steps(&self, run_id: Uuid) -> Result<(), EngineError> {
1775 interrupt_running_steps(self.store.as_ref(), run_id).await
1776 }
1777
1778 pub async fn fail_orphaned_steps(
1792 &self,
1793 run_id: Uuid,
1794 error_message: &str,
1795 ) -> Result<(), EngineError> {
1796 let steps = self.store.list_steps(run_id).await?;
1797 let now = Utc::now();
1798
1799 for step in steps {
1800 if step.status.state.is_terminal() {
1801 continue;
1802 }
1803
1804 let (target_status, error) = match step.status.state {
1805 StepStatus::Running | StepStatus::AwaitingApproval => {
1806 let err = if step.error.is_some() {
1807 None
1808 } else {
1809 Some(error_message.to_string())
1810 };
1811 (StepStatus::Failed, err)
1812 }
1813 StepStatus::Pending => (StepStatus::Skipped, None),
1814 _ => continue,
1815 };
1816
1817 if let Err(e) = self
1818 .store
1819 .update_step(
1820 step.id,
1821 StepUpdate {
1822 status: Some(target_status),
1823 error,
1824 completed_at: Some(now),
1825 ..StepUpdate::default()
1826 },
1827 )
1828 .await
1829 {
1830 warn!(
1831 run_id = %run_id,
1832 step_id = %step.id,
1833 step_name = %step.name,
1834 error = %e,
1835 "failed to cleanup orphaned step"
1836 );
1837 } else {
1838 info!(
1839 run_id = %run_id,
1840 step_id = %step.id,
1841 step_name = %step.name,
1842 from = %step.status.state,
1843 to = %target_status,
1844 "cleaned up orphaned step"
1845 );
1846 }
1847 }
1848
1849 Ok(())
1850 }
1851
1852 async fn release_then_execute(
1857 &self,
1858 run_id: Uuid,
1859 handler: &dyn WorkflowHandler,
1860 ctx: &mut WorkflowContext,
1861 ) -> Result<(), EngineError> {
1862 match self.provider.release_run(&run_id.to_string()).await {
1863 Ok(()) => handler.execute(ctx).await,
1864 Err(e) => Err(EngineError::Operation(OperationError::Agent(e))),
1865 }
1866 }
1867
1868 async fn finalize_run(
1874 &self,
1875 run_id: Uuid,
1876 workflow_name: &str,
1877 result: Result<(), EngineError>,
1878 ctx: &WorkflowContext,
1879 run_start: Instant,
1880 run_labels: HashMap<String, String>,
1881 ) -> Result<WorkflowResult, EngineError> {
1882 let total_duration = ctx.carried_duration_ms() + run_start.elapsed().as_millis() as u64;
1885 let completed_at = Utc::now();
1886
1887 let final_status;
1888 let final_run;
1889
1890 match result {
1891 Ok(()) => {
1892 final_status = if ctx.has_allowed_failure() {
1893 RunStatus::Warning
1894 } else {
1895 RunStatus::Completed
1896 };
1897 final_run = self
1898 .store
1899 .update_run_returning(
1900 run_id,
1901 RunUpdate {
1902 status: Some(final_status),
1903 cost_usd: Some(ctx.total_cost_usd()),
1904 duration_ms: Some(total_duration),
1905 completed_at: Some(completed_at),
1906 output: ctx.output().cloned(),
1907 ..RunUpdate::default()
1908 },
1909 )
1910 .await?;
1911
1912 info!(
1913 run_id = %run_id,
1914 status = %final_status,
1915 cost_usd = %ctx.total_cost_usd(),
1916 duration_ms = total_duration,
1917 "run completed"
1918 );
1919 }
1920 Err(EngineError::ApprovalRequired {
1921 run_id: approval_run_id,
1922 step_id,
1923 ref message,
1924 }) => {
1925 final_status = RunStatus::AwaitingApproval;
1926 final_run = self
1927 .store
1928 .update_run_returning(
1929 run_id,
1930 RunUpdate {
1931 status: Some(RunStatus::AwaitingApproval),
1932 cost_usd: Some(ctx.total_cost_usd()),
1933 duration_ms: Some(total_duration),
1934 ..RunUpdate::default()
1935 },
1936 )
1937 .await?;
1938
1939 info!(
1940 run_id = %approval_run_id,
1941 step_id = %step_id,
1942 message = %message,
1943 "run awaiting approval"
1944 );
1945
1946 self.publish_approval_requested(approval_run_id, step_id, message)
1947 .await?;
1948 }
1949 Err(EngineError::ChildSuspended {
1950 run_id: child_run_id,
1951 ref cause,
1952 }) => {
1953 final_status = cause.suspension_status();
1954 final_run = self
1958 .store
1959 .update_run_returning(
1960 run_id,
1961 RunUpdate {
1962 status: Some(final_status),
1963 cost_usd: Some(ctx.total_cost_usd()),
1964 duration_ms: Some(total_duration),
1965 ..RunUpdate::default()
1966 },
1967 )
1968 .await?;
1969
1970 let leaf = cause.suspension_leaf();
1971 info!(
1972 run_id = %run_id,
1973 child_run_id = %child_run_id,
1974 status = %final_status,
1975 cause = %leaf,
1976 "run suspended with its child run"
1977 );
1978
1979 match leaf {
1980 EngineError::ApprovalRequired {
1981 run_id: approval_run_id,
1982 step_id,
1983 message,
1984 } => {
1985 self.publish_approval_requested(*approval_run_id, *step_id, message)
1986 .await?;
1987 }
1988 EngineError::SignalWaiting {
1989 run_id: wait_run_id,
1990 step_id,
1991 step_name,
1992 name,
1993 key,
1994 deadline_at,
1995 } => {
1996 self.event_publisher
1997 .publish(Event::SignalAwaited(SignalAwaitedEvent {
1998 run_id: *wait_run_id,
1999 step_id: *step_id,
2000 step_name: step_name.clone(),
2001 name: name.clone(),
2002 key: key.clone(),
2003 deadline_at: *deadline_at,
2004 at: Utc::now(),
2005 }));
2006 }
2007 _ => {}
2010 }
2011 }
2012 Err(EngineError::HumanInputRequired {
2013 run_id: input_run_id,
2014 step_id,
2015 ref message,
2016 }) => {
2017 final_status = RunStatus::AwaitingApproval;
2018 final_run = self
2019 .store
2020 .update_run_returning(
2021 run_id,
2022 RunUpdate {
2023 status: Some(RunStatus::AwaitingApproval),
2024 cost_usd: Some(ctx.total_cost_usd()),
2025 duration_ms: Some(total_duration),
2026 ..RunUpdate::default()
2027 },
2028 )
2029 .await?;
2030
2031 info!(
2033 run_id = %input_run_id,
2034 step_id = %step_id,
2035 message = %message,
2036 "run awaiting human input"
2037 );
2038 }
2039 Err(EngineError::DelaySleeping {
2040 run_id: delay_run_id,
2041 step_id,
2042 wake_at,
2043 }) => {
2044 final_status = RunStatus::Sleeping;
2045 final_run = self
2046 .store
2047 .update_run_returning(
2048 run_id,
2049 RunUpdate {
2050 status: Some(RunStatus::Sleeping),
2051 cost_usd: Some(ctx.total_cost_usd()),
2052 duration_ms: Some(total_duration),
2053 scheduled_at: Some(wake_at),
2054 ..RunUpdate::default()
2055 },
2056 )
2057 .await?;
2058
2059 info!(
2060 run_id = %delay_run_id,
2061 step_id = %step_id,
2062 wake_at = %wake_at,
2063 "run sleeping until delay elapses"
2064 );
2065 }
2066 Err(EngineError::SignalWaiting {
2067 run_id: wait_run_id,
2068 step_id,
2069 ref step_name,
2070 ref name,
2071 ref key,
2072 deadline_at,
2073 }) => {
2074 final_status = RunStatus::Sleeping;
2075 let waiting = self
2079 .store
2080 .suspend_run_on_signal(run_id, step_id, deadline_at)
2081 .await?;
2082 final_run = self
2083 .store
2084 .update_run_returning(
2085 run_id,
2086 RunUpdate {
2087 cost_usd: Some(ctx.total_cost_usd()),
2088 duration_ms: Some(total_duration),
2089 ..RunUpdate::default()
2090 },
2091 )
2092 .await?;
2093
2094 if waiting {
2095 self.event_publisher
2096 .publish(Event::SignalAwaited(SignalAwaitedEvent {
2097 run_id: wait_run_id,
2098 step_id,
2099 step_name: step_name.clone(),
2100 name: name.clone(),
2101 key: key.clone(),
2102 deadline_at,
2103 at: Utc::now(),
2104 }));
2105 }
2106
2107 info!(
2108 run_id = %wait_run_id,
2109 step_id = %step_id,
2110 signal = %name,
2111 key = %key,
2112 deadline_at = %deadline_at,
2113 waiting,
2114 "run sleeping until a signal arrives"
2115 );
2116 }
2117 Err(err) => {
2118 let guardrail_stop = matches!(
2122 err,
2123 EngineError::RunBudgetExceeded { .. } | EngineError::WorkflowGuardRejected(_)
2124 );
2125
2126 final_status = if guardrail_stop {
2127 if let Err(store_err) = self
2128 .store
2129 .update_run(
2130 run_id,
2131 RunUpdate {
2132 status: Some(RunStatus::Cancelled),
2133 error: Some(err.to_string()),
2134 cost_usd: Some(ctx.total_cost_usd()),
2135 duration_ms: Some(total_duration),
2136 completed_at: Some(completed_at),
2137 output: ctx.output().cloned(),
2138 ..RunUpdate::default()
2139 },
2140 )
2141 .await
2142 {
2143 error!(run_id = %run_id, store_error = %store_err, "failed to persist run cancellation");
2144 }
2145 if let Err(cleanup_err) = self
2146 .fail_orphaned_steps(run_id, "run stopped: guardrail limit reached")
2147 .await
2148 {
2149 error!(run_id = %run_id, store_error = %cleanup_err, "failed to cleanup orphaned steps");
2150 }
2151 RunStatus::Cancelled
2152 } else {
2153 if let Some(output) = ctx.output()
2156 && let Err(store_err) = self
2157 .store
2158 .update_run(
2159 run_id,
2160 RunUpdate {
2161 output: Some(output.clone()),
2162 ..RunUpdate::default()
2163 },
2164 )
2165 .await
2166 {
2167 error!(run_id = %run_id, store_error = %store_err, "failed to persist run output");
2168 }
2169 self.fail_or_schedule_retry(
2170 run_id,
2171 &err.to_string(),
2172 is_run_retryable(&err),
2173 Some(ctx.total_cost_usd()),
2174 Some(total_duration),
2175 )
2176 .await
2177 .unwrap_or_else(|store_err| {
2178 error!(run_id = %run_id, store_error = %store_err, "failed to persist run failure");
2179 RunStatus::Failed
2180 })
2181 };
2182
2183 if matches!(err, EngineError::RunBudgetExceeded { .. }) {
2184 self.on_run_budget_exceeded(workflow_name, run_id, &err);
2185 }
2186
2187 error!(run_id = %run_id, status = %final_status, error = %err, "run stopped");
2188
2189 self.publish_run_status_changed(
2190 workflow_name,
2191 run_id,
2192 final_status,
2193 Some(err.to_string()),
2194 ctx,
2195 total_duration,
2196 run_labels,
2197 );
2198
2199 #[cfg(feature = "prometheus")]
2200 self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
2201
2202 return Err(err);
2203 }
2204 }
2205
2206 self.publish_run_status_changed(
2207 workflow_name,
2208 run_id,
2209 final_status,
2210 None,
2211 ctx,
2212 total_duration,
2213 run_labels,
2214 );
2215
2216 #[cfg(feature = "prometheus")]
2217 self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
2218
2219 Ok(WorkflowResult {
2220 run: final_run,
2221 steps: ctx.step_results().to_vec(),
2222 })
2223 }
2224
2225 async fn publish_approval_requested(
2228 &self,
2229 run_id: Uuid,
2230 step_id: Uuid,
2231 message: &str,
2232 ) -> Result<(), EngineError> {
2233 let requirement = self
2234 .store
2235 .get_step(step_id)
2236 .await?
2237 .and_then(|s| s.approval_requirement);
2238 self.event_publisher
2239 .publish(Event::ApprovalRequested(ApprovalRequestedEvent {
2240 run_id,
2241 step_id,
2242 message: message.to_string(),
2243 requirement,
2244 at: Utc::now(),
2245 }));
2246 Ok(())
2247 }
2248
2249 pub async fn fail_ancestors(&self, run_id: Uuid, reason: &str) -> Result<(), EngineError> {
2278 let mut current = self
2279 .store
2280 .get_run(run_id)
2281 .await?
2282 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
2283 let mut visited = HashSet::from([run_id]);
2285
2286 while let Some(parent_id) = chain_parent(¤t) {
2287 if !visited.insert(parent_id) {
2288 break;
2289 }
2290 let status = self
2291 .fail_or_schedule_retry(parent_id, reason, false, None, None)
2292 .await?;
2293 info!(
2294 run_id = %run_id,
2295 ancestor_run_id = %parent_id,
2296 status = %status,
2297 "ancestor run failed with its child"
2298 );
2299 current = self
2300 .store
2301 .get_run(parent_id)
2302 .await?
2303 .ok_or(EngineError::Store(StoreError::RunNotFound(parent_id)))?;
2304 }
2305
2306 Ok(())
2307 }
2308
2309 #[cfg(feature = "prometheus")]
2311 fn emit_run_metrics(
2312 &self,
2313 workflow_name: &str,
2314 status: RunStatus,
2315 duration_ms: u64,
2316 ctx: &WorkflowContext,
2317 ) {
2318 let status_str = status.to_string();
2319 let wf = workflow_name.to_string();
2320
2321 counter!(RUNS_TOTAL, "workflow" => wf.clone(), "status" => status_str.clone()).increment(1);
2322 histogram!(RUN_DURATION_SECONDS, "workflow" => wf.clone(), "status" => status_str)
2323 .record(duration_ms as f64 / 1000.0);
2324 histogram!(RUN_COST_USD, "workflow" => wf.clone()).record(
2325 ctx.total_cost_usd()
2326 .to_string()
2327 .parse::<f64>()
2328 .unwrap_or(0.0),
2329 );
2330 gauge!(RUNS_ACTIVE, "workflow" => wf).decrement(1.0);
2331 }
2332
2333 fn on_run_budget_exceeded(&self, workflow_name: &str, run_id: Uuid, err: &EngineError) {
2339 let EngineError::RunBudgetExceeded {
2340 limit_usd,
2341 spent_usd,
2342 step_budget_usd,
2343 ..
2344 } = err
2345 else {
2346 return;
2347 };
2348
2349 #[cfg(feature = "prometheus")]
2350 counter!(
2351 RUN_BUDGET_EXCEEDED_TOTAL,
2352 "workflow" => workflow_name.to_string(),
2353 "scope" => "run",
2354 )
2355 .increment(1);
2356
2357 self.event_publisher
2358 .publish(Event::RunBudgetExceeded(RunBudgetExceededEvent {
2359 run_id,
2360 workflow_name: workflow_name.to_string(),
2361 limit_usd: *limit_usd,
2362 spent_usd: *spent_usd,
2363 step_budget_usd: *step_budget_usd,
2364 at: Utc::now(),
2365 }));
2366 }
2367
2368 #[allow(clippy::too_many_arguments)]
2373 fn publish_run_status_changed(
2374 &self,
2375 workflow_name: &str,
2376 run_id: Uuid,
2377 to: RunStatus,
2378 error: Option<String>,
2379 ctx: &WorkflowContext,
2380 duration_ms: u64,
2381 labels: HashMap<String, String>,
2382 ) {
2383 let now = Utc::now();
2384 let cost_usd = ctx.total_cost_usd();
2385 let wf = workflow_name.to_string();
2386
2387 self.event_publisher
2388 .publish(Event::RunStatusChanged(RunStatusChangedEvent {
2389 run_id,
2390 workflow_name: wf.clone(),
2391 from: RunStatus::Running,
2392 to,
2393 error: error.clone(),
2394 cost_usd,
2395 duration_ms,
2396 labels: labels.clone(),
2397 at: now,
2398 }));
2399
2400 if to == RunStatus::Failed {
2401 self.event_publisher
2402 .publish(Event::RunFailed(RunFailedEvent {
2403 run_id,
2404 workflow_name: wf,
2405 error,
2406 cost_usd,
2407 duration_ms,
2408 labels,
2409 at: now,
2410 }));
2411 }
2412 }
2413}
2414
2415impl fmt::Debug for Engine {
2416 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
2417 f.debug_struct("Engine")
2418 .field("handlers", &self.handlers.keys().collect::<Vec<_>>())
2419 .finish_non_exhaustive()
2420 }
2421}
2422
2423#[cfg(test)]
2424mod tests {
2425 use super::*;
2426 use crate::config::ShellConfig;
2427 use crate::handler::{HandlerFuture, WorkflowHandler};
2428 use ironflow_core::providers::claude::ClaudeCodeProvider;
2429 use ironflow_core::providers::record_replay::RecordReplayProvider;
2430 use ironflow_store::memory::InMemoryStore;
2431 use ironflow_store::models::StepStatus;
2432 use serde_json::json;
2433
2434 struct EchoWorkflow;
2436
2437 impl WorkflowHandler for EchoWorkflow {
2438 fn name(&self) -> &str {
2439 "echo-workflow"
2440 }
2441
2442 fn describe(&self) -> WorkflowInfo {
2443 WorkflowInfo {
2444 description: "A simple workflow that echoes hello".to_string(),
2445 source_code: None,
2446 sub_workflows: Vec::new(),
2447 category: None,
2448 version: self.version().map(str::to_string),
2449 compatible_versions: Vec::new(),
2450 input_schema: None,
2451 default_labels: HashMap::new(),
2452 schedule: self.schedule().cloned(),
2453 default_max_cost_usd: self.default_max_cost_usd(),
2454 }
2455 }
2456
2457 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2458 Box::pin(async move {
2459 ctx.shell("greet", ShellConfig::new("echo hello")).await?;
2460 Ok(())
2461 })
2462 }
2463 }
2464
2465 struct FailingWorkflow;
2467
2468 impl WorkflowHandler for FailingWorkflow {
2469 fn name(&self) -> &str {
2470 "failing-workflow"
2471 }
2472
2473 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2474 Box::pin(async move {
2475 ctx.shell("fail", ShellConfig::new("exit 1")).await?;
2476 Ok(())
2477 })
2478 }
2479 }
2480
2481 fn create_test_engine() -> Engine {
2482 let store = Arc::new(InMemoryStore::new());
2483 let inner = ClaudeCodeProvider::new();
2484 let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
2485 inner,
2486 "/tmp/ironflow-fixtures",
2487 ));
2488 Engine::new(store, provider)
2489 }
2490
2491 #[test]
2492 fn engine_new_creates_instance() {
2493 let engine = create_test_engine();
2494 assert_eq!(engine.handler_names().len(), 0);
2495 }
2496
2497 #[test]
2498 fn execution_mode_defaults_to_local() {
2499 let engine = create_test_engine();
2500 assert_eq!(engine.execution_mode(), ExecutionMode::Local);
2501 }
2502
2503 #[test]
2504 fn with_execution_mode_overrides_the_default() {
2505 let engine = create_test_engine().with_execution_mode(ExecutionMode::Workers);
2506 assert_eq!(engine.execution_mode(), ExecutionMode::Workers);
2507 }
2508
2509 #[test]
2510 fn engine_register_handler() {
2511 let mut engine = create_test_engine();
2512 let result = engine.register(EchoWorkflow);
2513 assert!(result.is_ok());
2514 assert_eq!(engine.handler_names().len(), 1);
2515 assert!(engine.handler_names().contains(&"echo-workflow"));
2516 }
2517
2518 #[test]
2519 fn engine_register_duplicate_returns_error() {
2520 let mut engine = create_test_engine();
2521 engine.register(EchoWorkflow).unwrap();
2522 let result = engine.register(EchoWorkflow);
2523 assert!(result.is_err());
2524 }
2525
2526 #[test]
2527 fn engine_get_handler_found() {
2528 let mut engine = create_test_engine();
2529 engine.register(EchoWorkflow).unwrap();
2530 let handler = engine.get_handler("echo-workflow");
2531 assert!(handler.is_some());
2532 }
2533
2534 #[test]
2535 fn engine_get_handler_not_found() {
2536 let engine = create_test_engine();
2537 let handler = engine.get_handler("nonexistent");
2538 assert!(handler.is_none());
2539 }
2540
2541 #[test]
2542 fn engine_handler_names_lists_all() {
2543 let mut engine = create_test_engine();
2544 engine.register(EchoWorkflow).unwrap();
2545 engine.register(FailingWorkflow).unwrap();
2546 let names = engine.handler_names();
2547 assert_eq!(names.len(), 2);
2548 assert!(names.contains(&"echo-workflow"));
2549 assert!(names.contains(&"failing-workflow"));
2550 }
2551
2552 #[test]
2553 fn engine_handler_info_returns_description() {
2554 let mut engine = create_test_engine();
2555 engine.register(EchoWorkflow).unwrap();
2556 let info = engine.handler_info("echo-workflow");
2557 assert!(info.is_some());
2558 let info = info.unwrap();
2559 assert_eq!(info.description, "A simple workflow that echoes hello");
2560 }
2561
2562 struct CategorizedWorkflow;
2563
2564 impl WorkflowHandler for CategorizedWorkflow {
2565 fn name(&self) -> &str {
2566 "categorized"
2567 }
2568 fn category(&self) -> Option<&str> {
2569 Some("data/etl")
2570 }
2571 fn execute<'a>(
2572 &'a self,
2573 _ctx: &'a mut WorkflowContext,
2574 ) -> crate::handler::HandlerFuture<'a> {
2575 Box::pin(async move { Ok(()) })
2576 }
2577 }
2578
2579 #[test]
2580 fn engine_default_describe_propagates_category() {
2581 let mut engine = create_test_engine();
2582 engine.register(CategorizedWorkflow).unwrap();
2583 let info = engine.handler_info("categorized").unwrap();
2584 assert_eq!(info.category.as_deref(), Some("data/etl"));
2585 }
2586
2587 #[test]
2588 fn engine_default_describe_without_category() {
2589 let mut engine = create_test_engine();
2590 engine.register(EchoWorkflow).unwrap();
2591 let info = engine.handler_info("echo-workflow").unwrap();
2592 assert!(info.category.is_none());
2593 }
2594
2595 struct ScheduledWorkflow {
2600 schedule: CronSchedule,
2601 }
2602
2603 impl ScheduledWorkflow {
2604 fn new() -> Self {
2605 Self {
2606 schedule: CronSchedule::new("0 0 * * * *").unwrap(),
2607 }
2608 }
2609 }
2610
2611 impl WorkflowHandler for ScheduledWorkflow {
2612 fn name(&self) -> &str {
2613 "scheduled"
2614 }
2615 fn schedule(&self) -> Option<&CronSchedule> {
2616 Some(&self.schedule)
2617 }
2618 fn execute<'a>(
2619 &'a self,
2620 _ctx: &'a mut WorkflowContext,
2621 ) -> crate::handler::HandlerFuture<'a> {
2622 Box::pin(async move { Ok(()) })
2623 }
2624 }
2625
2626 #[test]
2627 fn engine_default_describe_propagates_schedule() {
2628 let mut engine = create_test_engine();
2629 engine.register(ScheduledWorkflow::new()).unwrap();
2630 let info = engine.handler_info("scheduled").unwrap();
2631 assert_eq!(
2632 info.schedule.as_ref().map(|s| s.as_str()),
2633 Some("0 0 * * * *")
2634 );
2635 }
2636
2637 #[test]
2638 fn engine_default_describe_without_schedule() {
2639 let mut engine = create_test_engine();
2640 engine.register(EchoWorkflow).unwrap();
2641 let info = engine.handler_info("echo-workflow").unwrap();
2642 assert!(info.schedule.is_none());
2643 }
2644
2645 #[test]
2646 fn scheduled_handlers_returns_only_scheduled() {
2647 let mut engine = create_test_engine();
2648 engine.register(EchoWorkflow).unwrap();
2649 engine.register(ScheduledWorkflow::new()).unwrap();
2650 engine.register(FailingWorkflow).unwrap();
2651
2652 let scheduled = engine.scheduled_handlers();
2653 assert_eq!(scheduled.len(), 1);
2654 assert_eq!(scheduled[0].0, "scheduled");
2655 assert_eq!(scheduled[0].1.as_str(), "0 0 * * * *");
2656 }
2657
2658 #[test]
2659 fn scheduled_handlers_empty_when_none_scheduled() {
2660 let mut engine = create_test_engine();
2661 engine.register(EchoWorkflow).unwrap();
2662 engine.register(FailingWorkflow).unwrap();
2663
2664 let scheduled = engine.scheduled_handlers();
2665 assert!(scheduled.is_empty());
2666 }
2667
2668 struct BadCategoryWorkflow(&'static str);
2669
2670 impl WorkflowHandler for BadCategoryWorkflow {
2671 fn name(&self) -> &str {
2672 "bad-category"
2673 }
2674 fn category(&self) -> Option<&str> {
2675 Some(self.0)
2676 }
2677 fn execute<'a>(
2678 &'a self,
2679 _ctx: &'a mut WorkflowContext,
2680 ) -> crate::handler::HandlerFuture<'a> {
2681 Box::pin(async move { Ok(()) })
2682 }
2683 }
2684
2685 #[test]
2686 fn engine_register_rejects_empty_category() {
2687 let mut engine = create_test_engine();
2688 let err = engine.register(BadCategoryWorkflow("")).unwrap_err();
2689 match err {
2690 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty category")),
2691 other => panic!("expected InvalidWorkflow, got {other:?}"),
2692 }
2693 }
2694
2695 #[test]
2696 fn engine_register_rejects_leading_slash_category() {
2697 let mut engine = create_test_engine();
2698 let err = engine
2699 .register(BadCategoryWorkflow("/data/etl"))
2700 .unwrap_err();
2701 match err {
2702 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("leading '/'")),
2703 other => panic!("expected InvalidWorkflow, got {other:?}"),
2704 }
2705 }
2706
2707 #[test]
2708 fn engine_register_rejects_trailing_slash_category() {
2709 let mut engine = create_test_engine();
2710 let err = engine
2711 .register(BadCategoryWorkflow("data/etl/"))
2712 .unwrap_err();
2713 match err {
2714 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("trailing '/'")),
2715 other => panic!("expected InvalidWorkflow, got {other:?}"),
2716 }
2717 }
2718
2719 #[test]
2720 fn engine_register_rejects_double_slash_category() {
2721 let mut engine = create_test_engine();
2722 let err = engine
2723 .register(BadCategoryWorkflow("data//etl"))
2724 .unwrap_err();
2725 match err {
2726 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty segment")),
2727 other => panic!("expected InvalidWorkflow, got {other:?}"),
2728 }
2729 }
2730
2731 #[test]
2732 fn engine_register_rejects_whitespace_only_segment_category() {
2733 let mut engine = create_test_engine();
2734 let err = engine
2735 .register(BadCategoryWorkflow("data/ /etl"))
2736 .unwrap_err();
2737 match err {
2738 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("whitespace-only segment")),
2739 other => panic!("expected InvalidWorkflow, got {other:?}"),
2740 }
2741 }
2742
2743 #[test]
2744 fn engine_register_accepts_valid_nested_category() {
2745 let mut engine = create_test_engine();
2746 assert!(engine.register(CategorizedWorkflow).is_ok());
2747 }
2748
2749 #[tokio::test]
2750 async fn engine_unknown_workflow_returns_error() {
2751 let engine = create_test_engine();
2752 let result = engine
2753 .run_handler("unknown", TriggerKind::Manual, json!({}))
2754 .await;
2755 assert!(result.is_err());
2756 match result {
2757 Err(EngineError::InvalidWorkflow(msg)) => {
2758 assert!(msg.contains("no handler registered"));
2759 }
2760 _ => panic!("expected InvalidWorkflow error"),
2761 }
2762 }
2763
2764 #[tokio::test]
2765 async fn engine_enqueue_handler_creates_pending_run() {
2766 let mut engine = create_test_engine();
2767 engine.register(EchoWorkflow).unwrap();
2768
2769 let run = engine
2770 .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
2771 .await
2772 .unwrap();
2773 assert_eq!(run.status.state, RunStatus::Pending);
2774 assert_eq!(run.workflow_name, "echo-workflow");
2775 }
2776
2777 #[tokio::test]
2778 async fn enqueue_handler_leaves_the_run_unattributed() {
2779 let mut engine = create_test_engine();
2780 engine.register(EchoWorkflow).unwrap();
2781
2782 let run = engine
2783 .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
2784 .await
2785 .unwrap();
2786
2787 assert!(run.created_by.is_none());
2788 }
2789
2790 #[tokio::test]
2791 async fn enqueue_handler_with_options_records_the_author() {
2792 let mut engine = create_test_engine();
2793 engine.register(EchoWorkflow).unwrap();
2794 let actor = RunActor::User {
2795 user_id: Uuid::now_v7(),
2796 };
2797
2798 let run = engine
2799 .enqueue_handler_with_options(
2800 "echo-workflow",
2801 TriggerKind::Api,
2802 json!({}),
2803 EnqueueOptions {
2804 created_by: Some(actor.clone()),
2805 ..Default::default()
2806 },
2807 )
2808 .await
2809 .unwrap()
2810 .into_run();
2811
2812 assert_eq!(run.created_by, Some(actor));
2813 }
2814
2815 #[tokio::test]
2816 async fn enqueue_handler_with_options_accepts_no_author() {
2817 let mut engine = create_test_engine();
2818 engine.register(EchoWorkflow).unwrap();
2819
2820 let run = engine
2821 .enqueue_handler_with_options(
2822 "echo-workflow",
2823 TriggerKind::Cron {
2824 schedule: "0 * * * * *".to_string(),
2825 },
2826 json!({}),
2827 EnqueueOptions::default(),
2828 )
2829 .await
2830 .unwrap()
2831 .into_run();
2832
2833 assert!(run.created_by.is_none());
2834 }
2835
2836 #[tokio::test]
2837 async fn enqueue_handler_with_options_stores_concurrency_limits() {
2838 let mut engine = create_test_engine();
2839 engine.register(EchoWorkflow).unwrap();
2840 let limits = vec![
2841 ConcurrencyLimit::new("repo:acme", 2),
2842 ConcurrencyLimit::new("tenant:42", 5),
2843 ];
2844
2845 let run = engine
2846 .enqueue_handler_with_options(
2847 "echo-workflow",
2848 TriggerKind::Api,
2849 json!({}),
2850 EnqueueOptions {
2851 concurrency_limits: limits.clone(),
2852 ..Default::default()
2853 },
2854 )
2855 .await
2856 .unwrap()
2857 .into_run();
2858
2859 assert_eq!(run.concurrency_limits, limits);
2860 }
2861
2862 #[tokio::test]
2863 async fn enqueue_rejects_invalid_concurrency_limits() {
2864 let mut engine = create_test_engine();
2865 engine.register(EchoWorkflow).unwrap();
2866
2867 let invalid = [
2868 vec![ConcurrencyLimit::new("repo:acme", 0)],
2869 vec![ConcurrencyLimit::new("", 1)],
2870 vec![
2871 ConcurrencyLimit::new("repo:acme", 1),
2872 ConcurrencyLimit::new("repo:acme", 2),
2873 ],
2874 ];
2875 for concurrency_limits in invalid {
2876 let err = engine
2877 .enqueue_handler_with_options(
2878 "echo-workflow",
2879 TriggerKind::Api,
2880 json!({}),
2881 EnqueueOptions {
2882 concurrency_limits,
2883 ..Default::default()
2884 },
2885 )
2886 .await
2887 .unwrap_err();
2888 assert!(
2889 matches!(err, EngineError::InvalidConcurrencyLimit(_)),
2890 "{err:?}"
2891 );
2892 }
2893
2894 let err = engine
2896 .enqueue_handler_with_options(
2897 "not-registered",
2898 TriggerKind::Api,
2899 json!({}),
2900 EnqueueOptions {
2901 concurrency_limits: vec![ConcurrencyLimit::new("repo:acme", 0)],
2902 ..Default::default()
2903 },
2904 )
2905 .await
2906 .unwrap_err();
2907 assert!(
2908 matches!(err, EngineError::InvalidConcurrencyLimit(_)),
2909 "{err:?}"
2910 );
2911
2912 let page = engine
2913 .store()
2914 .list_runs(RunFilter::default(), 1, 10)
2915 .await
2916 .unwrap();
2917 assert_eq!(page.total, 0, "no run may be created");
2918 }
2919
2920 #[tokio::test]
2921 async fn run_handler_leaves_the_run_unattributed() {
2922 let mut engine = create_test_engine();
2923 engine.register(EchoWorkflow).unwrap();
2924
2925 let run = engine
2926 .run_handler("echo-workflow", TriggerKind::Manual, json!({}))
2927 .await
2928 .unwrap()
2929 .run;
2930
2931 assert!(run.created_by.is_none());
2932 }
2933
2934 #[tokio::test]
2935 async fn engine_register_boxed() {
2936 let mut engine = create_test_engine();
2937 let handler: Box<dyn WorkflowHandler> = Box::new(EchoWorkflow);
2938 let result = engine.register_boxed(handler);
2939 assert!(result.is_ok());
2940 assert_eq!(engine.handler_names().len(), 1);
2941 }
2942
2943 #[tokio::test]
2944 async fn engine_store_and_provider_accessors() {
2945 let store = Arc::new(InMemoryStore::new());
2946 let inner = ClaudeCodeProvider::new();
2947 let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
2948 inner,
2949 "/tmp/ironflow-fixtures",
2950 ));
2951 let engine = Engine::new(store.clone(), provider.clone());
2952
2953 let _ = engine.store();
2955 let _ = engine.provider();
2956 }
2957
2958 use crate::operation::{Operation, OperationContext};
2963 use async_trait::async_trait;
2964 use ironflow_core::error::OperationError;
2965 use ironflow_store::models::StepKind;
2966
2967 struct FakeGitlabOp {
2968 project_id: u64,
2969 title: String,
2970 }
2971
2972 #[async_trait]
2973 impl Operation for FakeGitlabOp {
2974 fn kind(&self) -> &str {
2975 "gitlab"
2976 }
2977
2978 async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
2979 Ok(json!({
2980 "issue_id": 42,
2981 "project_id": self.project_id,
2982 "title": self.title,
2983 }))
2984 }
2985
2986 fn input(&self) -> Option<Value> {
2987 Some(json!({
2988 "project_id": self.project_id,
2989 "title": self.title,
2990 }))
2991 }
2992 }
2993
2994 struct FailingOp;
2995
2996 #[async_trait]
2997 impl Operation for FailingOp {
2998 fn kind(&self) -> &str {
2999 "broken-service"
3000 }
3001
3002 async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
3003 Err(OperationError::Http {
3004 status: None,
3005 message: "service unavailable".to_string(),
3006 })
3007 }
3008 }
3009
3010 struct OperationWorkflow;
3011
3012 impl WorkflowHandler for OperationWorkflow {
3013 fn name(&self) -> &str {
3014 "operation-workflow"
3015 }
3016
3017 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3018 Box::pin(async move {
3019 let op = FakeGitlabOp {
3020 project_id: 123,
3021 title: "Bug report".to_string(),
3022 };
3023 ctx.operation("create-issue", &op).await?;
3024 Ok(())
3025 })
3026 }
3027 }
3028
3029 struct FailingOperationWorkflow;
3030
3031 impl WorkflowHandler for FailingOperationWorkflow {
3032 fn name(&self) -> &str {
3033 "failing-operation-workflow"
3034 }
3035
3036 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3037 Box::pin(async move {
3038 ctx.operation("broken-call", &FailingOp).await?;
3039 Ok(())
3040 })
3041 }
3042 }
3043
3044 struct MixedWorkflow;
3045
3046 impl WorkflowHandler for MixedWorkflow {
3047 fn name(&self) -> &str {
3048 "mixed-workflow"
3049 }
3050
3051 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3052 Box::pin(async move {
3053 ctx.shell("build", ShellConfig::new("echo built")).await?;
3054 let op = FakeGitlabOp {
3055 project_id: 456,
3056 title: "Deploy done".to_string(),
3057 };
3058 let result = ctx.operation("notify-gitlab", &op).await?;
3059 assert_eq!(result.output["issue_id"], 42);
3060 Ok(())
3061 })
3062 }
3063 }
3064
3065 #[tokio::test]
3066 async fn operation_step_happy_path() {
3067 let mut engine = create_test_engine();
3068 engine.register(OperationWorkflow).unwrap();
3069
3070 let run = engine
3071 .run_handler("operation-workflow", TriggerKind::Manual, json!({}))
3072 .await
3073 .unwrap()
3074 .run;
3075
3076 assert_eq!(run.status.state, RunStatus::Completed);
3077
3078 let steps = engine.store().list_steps(run.id).await.unwrap();
3079
3080 assert_eq!(steps.len(), 1);
3081 assert_eq!(steps[0].name, "create-issue");
3082 assert_eq!(steps[0].kind, StepKind::Custom("gitlab".to_string()));
3083 assert_eq!(
3084 steps[0].status.state,
3085 ironflow_store::models::StepStatus::Completed
3086 );
3087
3088 let output = steps[0].output.as_ref().unwrap();
3089 assert_eq!(output["issue_id"], 42);
3090 assert_eq!(output["project_id"], 123);
3091
3092 let input = steps[0].input.as_ref().unwrap();
3093 assert_eq!(input["project_id"], 123);
3094 assert_eq!(input["title"], "Bug report");
3095 }
3096
3097 #[tokio::test]
3098 async fn operation_step_failure_marks_run_failed() {
3099 let mut engine = create_test_engine();
3100 engine.register(FailingOperationWorkflow).unwrap();
3101
3102 let result = engine
3103 .run_handler("failing-operation-workflow", TriggerKind::Manual, json!({}))
3104 .await;
3105
3106 assert!(result.is_err());
3107 }
3108
3109 #[tokio::test]
3110 async fn operation_mixed_with_shell_steps() {
3111 let mut engine = create_test_engine();
3112 engine.register(MixedWorkflow).unwrap();
3113
3114 let run = engine
3115 .run_handler("mixed-workflow", TriggerKind::Manual, json!({}))
3116 .await
3117 .unwrap()
3118 .run;
3119
3120 assert_eq!(run.status.state, RunStatus::Completed);
3121
3122 let steps = engine.store().list_steps(run.id).await.unwrap();
3123
3124 assert_eq!(steps.len(), 2);
3125 assert_eq!(steps[0].kind, StepKind::Shell);
3126 assert_eq!(steps[1].kind, StepKind::Custom("gitlab".to_string()));
3127 assert_eq!(steps[0].position, 0);
3128 assert_eq!(steps[1].position, 1);
3129 }
3130
3131 use crate::config::ApprovalConfig;
3136
3137 struct SingleApprovalWorkflow;
3138
3139 impl WorkflowHandler for SingleApprovalWorkflow {
3140 fn name(&self) -> &str {
3141 "single-approval"
3142 }
3143
3144 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3145 Box::pin(async move {
3146 ctx.shell("build", ShellConfig::new("echo built")).await?;
3147 ctx.approval("gate", ApprovalConfig::new("OK?")).await?;
3148 ctx.shell("deploy", ShellConfig::new("echo deployed"))
3149 .await?;
3150 Ok(())
3151 })
3152 }
3153 }
3154
3155 struct DoubleApprovalWorkflow;
3156
3157 impl WorkflowHandler for DoubleApprovalWorkflow {
3158 fn name(&self) -> &str {
3159 "double-approval"
3160 }
3161
3162 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3163 Box::pin(async move {
3164 ctx.shell("build", ShellConfig::new("echo built")).await?;
3165 ctx.approval("staging-gate", ApprovalConfig::new("Deploy staging?"))
3166 .await?;
3167 ctx.shell("deploy-staging", ShellConfig::new("echo staging"))
3168 .await?;
3169 ctx.approval("prod-gate", ApprovalConfig::new("Deploy prod?"))
3170 .await?;
3171 ctx.shell("deploy-prod", ShellConfig::new("echo prod"))
3172 .await?;
3173 Ok(())
3174 })
3175 }
3176 }
3177
3178 #[tokio::test]
3179 async fn approval_pauses_run() {
3180 let mut engine = create_test_engine();
3181 engine.register(SingleApprovalWorkflow).unwrap();
3182
3183 let run = engine
3184 .run_handler("single-approval", TriggerKind::Manual, json!({}))
3185 .await
3186 .unwrap()
3187 .run;
3188
3189 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3190
3191 let steps = engine.store().list_steps(run.id).await.unwrap();
3192 assert_eq!(steps.len(), 2); assert_eq!(steps[0].kind, StepKind::Shell);
3194 assert_eq!(steps[0].status.state, StepStatus::Completed);
3195 assert_eq!(steps[1].kind, StepKind::Approval);
3196 assert_eq!(steps[1].status.state, StepStatus::AwaitingApproval);
3197 }
3198
3199 #[tokio::test]
3200 async fn approval_resume_completes_run() {
3201 let mut engine = create_test_engine();
3202 engine.register(SingleApprovalWorkflow).unwrap();
3203
3204 let run = engine
3206 .run_handler("single-approval", TriggerKind::Manual, json!({}))
3207 .await
3208 .unwrap()
3209 .run;
3210 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3211
3212 engine
3214 .store()
3215 .update_run_status(run.id, RunStatus::Running)
3216 .await
3217 .unwrap();
3218
3219 let resumed = engine.resume_run(run.id).await.unwrap().run;
3221 assert_eq!(resumed.status.state, RunStatus::Completed);
3222
3223 let steps = engine.store().list_steps(run.id).await.unwrap();
3224 assert_eq!(steps.len(), 3); assert_eq!(steps[0].name, "build");
3226 assert_eq!(steps[0].status.state, StepStatus::Completed);
3227 assert_eq!(steps[1].name, "gate");
3228 assert_eq!(steps[1].kind, StepKind::Approval);
3229 assert_eq!(steps[1].status.state, StepStatus::Completed);
3230 assert_eq!(steps[2].name, "deploy");
3231 assert_eq!(steps[2].status.state, StepStatus::Completed);
3232 }
3233
3234 #[tokio::test]
3235 async fn double_approval_two_resumes() {
3236 let mut engine = create_test_engine();
3237 engine.register(DoubleApprovalWorkflow).unwrap();
3238
3239 let run = engine
3241 .run_handler("double-approval", TriggerKind::Manual, json!({}))
3242 .await
3243 .unwrap()
3244 .run;
3245 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3246
3247 let steps = engine.store().list_steps(run.id).await.unwrap();
3248 assert_eq!(steps.len(), 2); engine
3252 .store()
3253 .update_run_status(run.id, RunStatus::Running)
3254 .await
3255 .unwrap();
3256
3257 let resumed = engine.resume_run(run.id).await.unwrap().run;
3258 assert_eq!(resumed.status.state, RunStatus::AwaitingApproval);
3259
3260 let steps = engine.store().list_steps(run.id).await.unwrap();
3261 assert_eq!(steps.len(), 4); engine
3265 .store()
3266 .update_run_status(run.id, RunStatus::Running)
3267 .await
3268 .unwrap();
3269
3270 let final_run = engine.resume_run(run.id).await.unwrap().run;
3271 assert_eq!(final_run.status.state, RunStatus::Completed);
3272
3273 let steps = engine.store().list_steps(run.id).await.unwrap();
3274 assert_eq!(steps.len(), 5);
3275 assert_eq!(steps[0].name, "build");
3276 assert_eq!(steps[1].name, "staging-gate");
3277 assert_eq!(steps[2].name, "deploy-staging");
3278 assert_eq!(steps[3].name, "prod-gate");
3279 assert_eq!(steps[4].name, "deploy-prod");
3280
3281 for step in &steps {
3282 assert_eq!(step.status.state, StepStatus::Completed);
3283 }
3284 }
3285
3286 use ironflow_store::models::{NewStep, StepUpdate, step_trace_id};
3291
3292 async fn create_step_with_status(
3293 store: &Arc<dyn Store>,
3294 run_id: Uuid,
3295 name: &str,
3296 position: u32,
3297 status: StepStatus,
3298 ) -> ironflow_store::models::Step {
3299 let step = store
3300 .create_step(NewStep {
3301 run_id,
3302 trace_id: step_trace_id(run_id, name, position),
3303 name: name.to_string(),
3304 kind: StepKind::Shell,
3305 position,
3306 input: None,
3307 is_error_handler: false,
3308 })
3309 .await
3310 .unwrap();
3311
3312 match status {
3313 StepStatus::Pending => {}
3314 StepStatus::Running => {
3315 store
3316 .update_step(
3317 step.id,
3318 StepUpdate {
3319 status: Some(StepStatus::Running),
3320 ..StepUpdate::default()
3321 },
3322 )
3323 .await
3324 .unwrap();
3325 }
3326 StepStatus::Completed => {
3327 store
3328 .update_step(
3329 step.id,
3330 StepUpdate {
3331 status: Some(StepStatus::Running),
3332 ..StepUpdate::default()
3333 },
3334 )
3335 .await
3336 .unwrap();
3337 store
3338 .update_step(
3339 step.id,
3340 StepUpdate {
3341 status: Some(StepStatus::Completed),
3342 ..StepUpdate::default()
3343 },
3344 )
3345 .await
3346 .unwrap();
3347 }
3348 StepStatus::AwaitingApproval => {
3349 store
3350 .update_step(
3351 step.id,
3352 StepUpdate {
3353 status: Some(StepStatus::Running),
3354 ..StepUpdate::default()
3355 },
3356 )
3357 .await
3358 .unwrap();
3359 store
3360 .update_step(
3361 step.id,
3362 StepUpdate {
3363 status: Some(StepStatus::AwaitingApproval),
3364 ..StepUpdate::default()
3365 },
3366 )
3367 .await
3368 .unwrap();
3369 }
3370 _ => panic!("unsupported status for test helper: {status}"),
3371 }
3372
3373 store.get_step(step.id).await.unwrap().unwrap()
3374 }
3375
3376 #[tokio::test]
3377 async fn fail_orphaned_steps_marks_running_as_failed() {
3378 let engine = create_test_engine();
3379 let run = engine
3380 .store()
3381 .create_run(NewRun {
3382 created_by: None,
3383 workflow_name: "test".to_string(),
3384 trigger: TriggerKind::Manual,
3385 payload: json!({}),
3386 max_retries: 0,
3387 handler_version: None,
3388 labels: HashMap::new(),
3389 scheduled_at: None,
3390 idempotency_key: None,
3391 concurrency_key: None,
3392 concurrency_limits: Vec::new(),
3393 max_cost_usd: None,
3394 })
3395 .await
3396 .unwrap()
3397 .into_run();
3398
3399 let step = create_step_with_status(
3400 engine.store(),
3401 run.id,
3402 "running-step",
3403 0,
3404 StepStatus::Running,
3405 )
3406 .await;
3407
3408 engine
3409 .fail_orphaned_steps(run.id, "parent run timed out")
3410 .await
3411 .unwrap();
3412
3413 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
3414 assert_eq!(updated.status.state, StepStatus::Failed);
3415 assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
3416 assert!(updated.completed_at.is_some());
3417 }
3418
3419 #[tokio::test]
3420 async fn fail_orphaned_steps_marks_pending_as_skipped() {
3421 let engine = create_test_engine();
3422 let run = engine
3423 .store()
3424 .create_run(NewRun {
3425 created_by: None,
3426 workflow_name: "test".to_string(),
3427 trigger: TriggerKind::Manual,
3428 payload: json!({}),
3429 max_retries: 0,
3430 handler_version: None,
3431 labels: HashMap::new(),
3432 scheduled_at: None,
3433 idempotency_key: None,
3434 concurrency_key: None,
3435 concurrency_limits: Vec::new(),
3436 max_cost_usd: None,
3437 })
3438 .await
3439 .unwrap()
3440 .into_run();
3441
3442 let step = create_step_with_status(
3443 engine.store(),
3444 run.id,
3445 "pending-step",
3446 0,
3447 StepStatus::Pending,
3448 )
3449 .await;
3450
3451 engine
3452 .fail_orphaned_steps(run.id, "parent run timed out")
3453 .await
3454 .unwrap();
3455
3456 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
3457 assert_eq!(updated.status.state, StepStatus::Skipped);
3458 assert!(updated.error.is_none());
3459 assert!(updated.completed_at.is_some());
3460 }
3461
3462 #[tokio::test]
3463 async fn fail_orphaned_steps_marks_awaiting_approval_as_failed() {
3464 let engine = create_test_engine();
3465 let run = engine
3466 .store()
3467 .create_run(NewRun {
3468 created_by: None,
3469 workflow_name: "test".to_string(),
3470 trigger: TriggerKind::Manual,
3471 payload: json!({}),
3472 max_retries: 0,
3473 handler_version: None,
3474 labels: HashMap::new(),
3475 scheduled_at: None,
3476 idempotency_key: None,
3477 concurrency_key: None,
3478 concurrency_limits: Vec::new(),
3479 max_cost_usd: None,
3480 })
3481 .await
3482 .unwrap()
3483 .into_run();
3484
3485 let step = create_step_with_status(
3486 engine.store(),
3487 run.id,
3488 "approval-step",
3489 0,
3490 StepStatus::AwaitingApproval,
3491 )
3492 .await;
3493
3494 engine
3495 .fail_orphaned_steps(run.id, "parent run timed out")
3496 .await
3497 .unwrap();
3498
3499 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
3500 assert_eq!(updated.status.state, StepStatus::Failed);
3501 assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
3502 assert!(updated.completed_at.is_some());
3503 }
3504
3505 #[tokio::test]
3506 async fn fail_orphaned_steps_skips_terminal_steps() {
3507 let engine = create_test_engine();
3508 let run = engine
3509 .store()
3510 .create_run(NewRun {
3511 created_by: None,
3512 workflow_name: "test".to_string(),
3513 trigger: TriggerKind::Manual,
3514 payload: json!({}),
3515 max_retries: 0,
3516 handler_version: None,
3517 labels: HashMap::new(),
3518 scheduled_at: None,
3519 idempotency_key: None,
3520 concurrency_key: None,
3521 concurrency_limits: Vec::new(),
3522 max_cost_usd: None,
3523 })
3524 .await
3525 .unwrap()
3526 .into_run();
3527
3528 let completed_step =
3529 create_step_with_status(engine.store(), run.id, "done", 0, StepStatus::Completed).await;
3530 let running_step =
3531 create_step_with_status(engine.store(), run.id, "in-flight", 1, StepStatus::Running)
3532 .await;
3533
3534 engine
3535 .fail_orphaned_steps(run.id, "parent run timed out")
3536 .await
3537 .unwrap();
3538
3539 let completed = engine
3540 .store()
3541 .get_step(completed_step.id)
3542 .await
3543 .unwrap()
3544 .unwrap();
3545 assert_eq!(completed.status.state, StepStatus::Completed);
3546
3547 let failed = engine
3548 .store()
3549 .get_step(running_step.id)
3550 .await
3551 .unwrap()
3552 .unwrap();
3553 assert_eq!(failed.status.state, StepStatus::Failed);
3554 }
3555
3556 #[tokio::test]
3557 async fn fail_orphaned_steps_mixed_states() {
3558 let engine = create_test_engine();
3559 let run = engine
3560 .store()
3561 .create_run(NewRun {
3562 created_by: None,
3563 workflow_name: "test".to_string(),
3564 trigger: TriggerKind::Manual,
3565 payload: json!({}),
3566 max_retries: 0,
3567 handler_version: None,
3568 labels: HashMap::new(),
3569 scheduled_at: None,
3570 idempotency_key: None,
3571 concurrency_key: None,
3572 concurrency_limits: Vec::new(),
3573 max_cost_usd: None,
3574 })
3575 .await
3576 .unwrap()
3577 .into_run();
3578
3579 let s_completed =
3580 create_step_with_status(engine.store(), run.id, "step-1", 0, StepStatus::Completed)
3581 .await;
3582 let s_running =
3583 create_step_with_status(engine.store(), run.id, "step-2", 1, StepStatus::Running).await;
3584 let s_pending =
3585 create_step_with_status(engine.store(), run.id, "step-3", 2, StepStatus::Pending).await;
3586
3587 engine.fail_orphaned_steps(run.id, "timeout").await.unwrap();
3588
3589 let r_completed = engine
3590 .store()
3591 .get_step(s_completed.id)
3592 .await
3593 .unwrap()
3594 .unwrap();
3595 assert_eq!(r_completed.status.state, StepStatus::Completed);
3596
3597 let r_running = engine
3598 .store()
3599 .get_step(s_running.id)
3600 .await
3601 .unwrap()
3602 .unwrap();
3603 assert_eq!(r_running.status.state, StepStatus::Failed);
3604 assert_eq!(r_running.error.as_deref(), Some("timeout"));
3605
3606 let r_pending = engine
3607 .store()
3608 .get_step(s_pending.id)
3609 .await
3610 .unwrap()
3611 .unwrap();
3612 assert_eq!(r_pending.status.state, StepStatus::Skipped);
3613 assert!(r_pending.error.is_none());
3614 }
3615
3616 #[tokio::test]
3617 async fn fail_orphaned_steps_no_steps_is_noop() {
3618 let engine = create_test_engine();
3619 let run = engine
3620 .store()
3621 .create_run(NewRun {
3622 created_by: None,
3623 workflow_name: "test".to_string(),
3624 trigger: TriggerKind::Manual,
3625 payload: json!({}),
3626 max_retries: 0,
3627 handler_version: None,
3628 labels: HashMap::new(),
3629 scheduled_at: None,
3630 idempotency_key: None,
3631 concurrency_key: None,
3632 concurrency_limits: Vec::new(),
3633 max_cost_usd: None,
3634 })
3635 .await
3636 .unwrap()
3637 .into_run();
3638
3639 let result = engine.fail_orphaned_steps(run.id, "timeout").await;
3640 assert!(result.is_ok());
3641 }
3642
3643 #[tokio::test]
3644 async fn fail_orphaned_steps_preserves_existing_error() {
3645 let engine = create_test_engine();
3646 let run = engine
3647 .store()
3648 .create_run(NewRun {
3649 created_by: None,
3650 workflow_name: "test".to_string(),
3651 trigger: TriggerKind::Manual,
3652 payload: json!({}),
3653 max_retries: 0,
3654 handler_version: None,
3655 labels: HashMap::new(),
3656 scheduled_at: None,
3657 idempotency_key: None,
3658 concurrency_key: None,
3659 concurrency_limits: Vec::new(),
3660 max_cost_usd: None,
3661 })
3662 .await
3663 .unwrap()
3664 .into_run();
3665
3666 let step_with_error = create_step_with_status(
3667 engine.store(),
3668 run.id,
3669 "already-errored",
3670 0,
3671 StepStatus::Running,
3672 )
3673 .await;
3674
3675 engine
3676 .store()
3677 .update_step(
3678 step_with_error.id,
3679 StepUpdate {
3680 error: Some("real error from provider".to_string()),
3681 ..StepUpdate::default()
3682 },
3683 )
3684 .await
3685 .unwrap();
3686
3687 let step_no_error = create_step_with_status(
3688 engine.store(),
3689 run.id,
3690 "no-error-yet",
3691 1,
3692 StepStatus::Running,
3693 )
3694 .await;
3695
3696 engine
3697 .fail_orphaned_steps(run.id, "parent run failed")
3698 .await
3699 .unwrap();
3700
3701 let updated_with = engine
3702 .store()
3703 .get_step(step_with_error.id)
3704 .await
3705 .unwrap()
3706 .unwrap();
3707 assert_eq!(updated_with.status.state, StepStatus::Failed);
3708 assert_eq!(
3709 updated_with.error.as_deref(),
3710 Some("real error from provider"),
3711 );
3712
3713 let updated_without = engine
3714 .store()
3715 .get_step(step_no_error.id)
3716 .await
3717 .unwrap()
3718 .unwrap();
3719 assert_eq!(updated_without.status.state, StepStatus::Failed);
3720 assert_eq!(updated_without.error.as_deref(), Some("parent run failed"),);
3721 }
3722}