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 NewRun, NewSignal, Run, RunActor, RunCreation, RunFilter, RunStatus, RunUpdate, SignalInsert,
30 SignalStepResolution, StepStatus, StepUpdate, TriggerKind,
31};
32use ironflow_store::store::Store;
33#[cfg(feature = "prometheus")]
34use metrics::{counter, gauge, histogram};
35
36use crate::artifact::ArtifactSink;
37use crate::budget::{BudgetConfig, month_start};
38use crate::context::{PARENT_RUN_ID_LABEL, WorkflowContext};
39use crate::error::EngineError;
40use crate::executor::{StepInterceptor, StepResult};
41use crate::guard::{WorkflowGuardConfig, new_shared_guard_state};
42use crate::handler::{WorkflowHandler, WorkflowInfo};
43use crate::log_sender::LogSender;
44use crate::notify::{
45 ApprovalRequestedEvent, Event, EventPublisher, EventSubscriber, RunBudgetExceededEvent,
46 RunFailedEvent, RunStatusChangedEvent, SignalAwaitedEvent, SignalReceivedEvent,
47 WorkflowEventBus,
48};
49use crate::plan::{
50 ExecutionPlan, PlanOptions, PlanRecorder, SharedPlanRecorder, estimate_durations, lock_plan,
51};
52use crate::retry_policy::{backoff_for_retry, is_run_retryable};
53use crate::schedule::CronSchedule;
54use crate::signal::{
55 Signal, SignalDelivery, SignalRejected, SignalResumed, received_output, validate_step_payload,
56};
57use ironflow_core::decision::DecisionProvider;
58
59#[derive(Debug, Clone)]
78pub struct WorkflowResult {
79 pub run: Run,
81 pub steps: Vec<StepResult>,
83}
84
85#[derive(Debug, Clone, Default)]
104pub struct EnqueueOptions {
105 pub max_retries: u32,
107 pub labels: HashMap<String, String>,
109 pub scheduled_at: Option<DateTime<Utc>>,
112 pub max_cost_usd: Option<Decimal>,
116 pub created_by: Option<RunActor>,
119 pub idempotency_key: Option<String>,
125}
126
127#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
138pub enum ExecutionMode {
139 #[default]
143 Local,
144 Workers,
147}
148
149pub struct Engine {
190 store: Arc<dyn Store>,
191 provider: Arc<dyn AgentProvider>,
192 handlers: HashMap<String, Arc<dyn WorkflowHandler>>,
193 event_publisher: EventPublisher,
194 log_sender: Option<LogSender>,
195 budget: BudgetConfig,
196 artifact_sink: Option<Arc<dyn ArtifactSink>>,
197 guard_config: Option<WorkflowGuardConfig>,
198 event_bus: Option<WorkflowEventBus>,
199 decision_provider: Option<Arc<dyn DecisionProvider>>,
200 step_interceptor: Option<Arc<dyn StepInterceptor>>,
201 execution_mode: ExecutionMode,
202}
203
204fn validate_category(handler_name: &str, category: &str) -> Result<(), EngineError> {
214 let reject = |reason: &str| {
215 Err(EngineError::InvalidWorkflow(format!(
216 "handler '{handler_name}' has invalid category '{category}': {reason}"
217 )))
218 };
219
220 if category.is_empty() {
221 return reject("empty category");
222 }
223 if category.starts_with('/') {
224 return reject("leading '/'");
225 }
226 if category.ends_with('/') {
227 return reject("trailing '/'");
228 }
229 for segment in category.split('/') {
230 if segment.is_empty() {
231 return reject("empty segment (double '/')");
232 }
233 if segment.trim().is_empty() {
234 return reject("whitespace-only segment");
235 }
236 }
237 Ok(())
238}
239
240fn chain_label(run: &Run, key: &str) -> Option<Uuid> {
245 if !matches!(run.trigger, TriggerKind::Workflow) {
246 return None;
247 }
248 let id = Uuid::parse_str(run.labels.get(key)?).ok()?;
249 (id != run.id).then_some(id)
250}
251
252pub(crate) fn chain_root(run: &Run) -> Option<Uuid> {
257 chain_label(run, LABEL_ROOT_RUN_ID)
258}
259
260fn chain_parent(run: &Run) -> Option<Uuid> {
262 chain_label(run, PARENT_RUN_ID_LABEL)
263}
264
265impl Engine {
266 pub fn new(store: Arc<dyn Store>, provider: Arc<dyn AgentProvider>) -> Self {
282 Self {
283 store,
284 provider,
285 handlers: HashMap::new(),
286 event_publisher: EventPublisher::new(),
287 log_sender: None,
288 budget: BudgetConfig::new(),
289 artifact_sink: None,
290 guard_config: None,
291 event_bus: None,
292 decision_provider: None,
293 step_interceptor: None,
294 execution_mode: ExecutionMode::default(),
295 }
296 }
297
298 pub fn with_decision_provider(mut self, provider: Arc<dyn DecisionProvider>) -> Self {
319 self.decision_provider = Some(provider);
320 self
321 }
322
323 pub fn with_step_interceptor(mut self, interceptor: Arc<dyn StepInterceptor>) -> Self {
347 self.step_interceptor = Some(interceptor);
348 self
349 }
350
351 pub fn step_interceptor(&self) -> Option<&Arc<dyn StepInterceptor>> {
353 self.step_interceptor.as_ref()
354 }
355
356 pub fn with_budget_config(mut self, budget: BudgetConfig) -> Self {
377 self.budget = budget;
378 self
379 }
380
381 pub fn budget_config(&self) -> &BudgetConfig {
383 &self.budget
384 }
385
386 pub fn with_guard_config(mut self, config: WorkflowGuardConfig) -> Self {
408 self.guard_config = Some(config);
409 self
410 }
411
412 pub fn guard_config(&self) -> Option<&WorkflowGuardConfig> {
414 self.guard_config.as_ref()
415 }
416
417 pub fn with_execution_mode(mut self, mode: ExecutionMode) -> Self {
438 self.execution_mode = mode;
439 self
440 }
441
442 pub fn execution_mode(&self) -> ExecutionMode {
444 self.execution_mode
445 }
446
447 pub fn set_log_sender(&mut self, sender: LogSender) {
453 self.log_sender = Some(sender);
454 }
455
456 pub fn set_artifact_sink(&mut self, sink: Arc<dyn ArtifactSink>) {
474 self.artifact_sink = Some(sink);
475 }
476
477 pub fn artifact_sink(&self) -> Option<&Arc<dyn ArtifactSink>> {
479 self.artifact_sink.as_ref()
480 }
481
482 pub fn set_event_bus(&mut self, bus: WorkflowEventBus) {
499 self.event_bus = Some(bus);
500 }
501
502 pub fn event_bus(&self) -> Option<&WorkflowEventBus> {
504 self.event_bus.as_ref()
505 }
506
507 pub fn store(&self) -> &Arc<dyn Store> {
509 &self.store
510 }
511
512 pub fn provider(&self) -> &Arc<dyn AgentProvider> {
514 &self.provider
515 }
516
517 fn build_context(&self, run: &Run) -> WorkflowContext {
526 let handlers = self.handlers.clone();
527 let resolver: crate::context::HandlerResolver =
528 Arc::new(move |name: &str| handlers.get(name).cloned());
529 let mut ctx = WorkflowContext::with_handler_resolver(
530 run.id,
531 run.workflow_name.clone(),
532 self.store.clone(),
533 self.provider.clone(),
534 resolver,
535 );
536 ctx.carry_over_run_totals(run.retry_count + 1, run.cost_usd, run.duration_ms);
537 ctx.set_max_cost_usd(run.max_cost_usd);
538 ctx.set_run_created_at(run.created_at);
539 if let Some(ref sender) = self.log_sender {
540 ctx.set_log_sender(sender.clone());
541 }
542 if let Some(ref sink) = self.artifact_sink {
543 ctx.set_artifact_sink(sink.clone());
544 }
545 if let Some(ref bus) = self.event_bus {
546 ctx.set_event_bus(bus.clone());
547 }
548 if let Some(ref provider) = self.decision_provider {
549 ctx.set_decision_provider(provider.clone());
550 }
551 if let Some(ref interceptor) = self.step_interceptor {
552 ctx.set_step_interceptor(interceptor.clone());
553 }
554 ctx
555 }
556
557 fn build_context_with_guard(
563 &self,
564 run: &Run,
565 handler: &dyn WorkflowHandler,
566 ) -> WorkflowContext {
567 let mut ctx = self.build_context(run);
568 let guard_config = handler.guard_config().or_else(|| self.guard_config.clone());
569 if let Some(config) = guard_config {
570 ctx.set_guard(config, new_shared_guard_state());
571 }
572 ctx
573 }
574
575 async fn check_monthly_quota(&self, workflow_name: &str) -> Result<(), EngineError> {
586 let Some(limit) = self.budget.monthly_cost_limit_usd else {
587 return Ok(());
588 };
589
590 let stats = self
591 .store
592 .get_stats(RunFilter {
593 created_after: Some(month_start(Utc::now())),
594 ..RunFilter::default()
595 })
596 .await?;
597
598 if stats.total_cost_usd < limit {
599 return Ok(());
600 }
601
602 warn!(
603 workflow = %workflow_name,
604 limit_usd = %limit,
605 spent_usd = %stats.total_cost_usd,
606 "monthly cost quota exhausted, refusing new run"
607 );
608
609 #[cfg(feature = "prometheus")]
610 counter!(
611 RUN_BUDGET_EXCEEDED_TOTAL,
612 "workflow" => workflow_name.to_string(),
613 "scope" => "monthly",
614 )
615 .increment(1);
616
617 Err(EngineError::MonthlyBudgetExceeded {
618 limit_usd: limit,
619 spent_usd: stats.total_cost_usd,
620 })
621 }
622
623 pub fn register(&mut self, handler: impl WorkflowHandler + 'static) -> Result<(), EngineError> {
667 let name = handler.name().to_string();
668 if self.handlers.contains_key(&name) {
669 return Err(EngineError::InvalidWorkflow(format!(
670 "handler '{}' already registered",
671 name
672 )));
673 }
674 if let Some(category) = handler.category() {
675 validate_category(&name, category)?;
676 }
677 self.handlers.insert(name, Arc::new(handler));
678 Ok(())
679 }
680
681 pub fn register_boxed(&mut self, handler: Box<dyn WorkflowHandler>) -> Result<(), EngineError> {
688 let name = handler.name().to_string();
689 if self.handlers.contains_key(&name) {
690 return Err(EngineError::InvalidWorkflow(format!(
691 "handler '{}' already registered",
692 name
693 )));
694 }
695 if let Some(category) = handler.category() {
696 validate_category(&name, category)?;
697 }
698 self.handlers.insert(name, Arc::from(handler));
699 Ok(())
700 }
701
702 pub fn get_handler(&self, name: &str) -> Option<&Arc<dyn WorkflowHandler>> {
704 self.handlers.get(name)
705 }
706
707 pub fn handler_names(&self) -> Vec<&str> {
709 self.handlers.keys().map(|s| s.as_str()).collect()
710 }
711
712 pub fn handler_info(&self, name: &str) -> Option<WorkflowInfo> {
714 self.handlers.get(name).map(|h| h.describe())
715 }
716
717 pub fn scheduled_handlers(&self) -> Vec<(&str, &CronSchedule)> {
741 self.handlers
742 .iter()
743 .filter_map(|(name, handler)| handler.schedule().map(|sched| (name.as_str(), sched)))
744 .collect()
745 }
746
747 pub fn subscribe(
772 &mut self,
773 subscriber: impl EventSubscriber + 'static,
774 event_types: &[&'static str],
775 ) {
776 self.event_publisher.subscribe(subscriber, event_types);
777 }
778
779 pub fn event_publisher(&self) -> &EventPublisher {
784 &self.event_publisher
785 }
786
787 #[tracing::instrument(name = "engine.run_handler", skip_all, fields(workflow = %handler_name))]
817 pub async fn run_handler(
818 &self,
819 handler_name: &str,
820 trigger: TriggerKind,
821 payload: Value,
822 ) -> Result<WorkflowResult, EngineError> {
823 let handler = self
824 .handlers
825 .get(handler_name)
826 .ok_or_else(|| {
827 EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
828 })?
829 .clone();
830
831 self.check_monthly_quota(handler_name).await?;
832
833 let handler_version = handler.version().map(str::to_string);
834 let max_cost_usd = self
835 .budget
836 .resolve_run_cap(None, handler.default_max_cost_usd());
837 let run = self
838 .store
839 .create_run(NewRun {
840 created_by: None,
841 workflow_name: handler_name.to_string(),
842 trigger,
843 payload,
844 max_retries: 0,
845 handler_version,
846 labels: handler.default_labels(),
847 scheduled_at: None,
848 idempotency_key: None,
849 max_cost_usd,
850 })
851 .await?
852 .into_run();
853
854 let run_id = run.id;
855 info!(run_id = %run_id, handler_version = run.handler_version.as_deref().unwrap_or(""), "run created");
856
857 self.store
858 .update_run_status(run_id, RunStatus::Running)
859 .await?;
860
861 #[cfg(feature = "prometheus")]
862 gauge!(RUNS_ACTIVE, "workflow" => handler_name.to_string()).increment(1.0);
863
864 let run_start = Instant::now();
865 let mut ctx = self.build_context_with_guard(&run, handler.as_ref());
866
867 let result = handler.execute(&mut ctx).await;
868 self.finalize_run(run_id, handler_name, result, &ctx, run_start, run.labels)
869 .await
870 }
871
872 #[tracing::instrument(name = "engine.plan_handler", skip_all, fields(workflow = %handler_name))]
914 pub async fn plan_handler(
915 &self,
916 handler_name: &str,
917 payload: Value,
918 options: PlanOptions,
919 ) -> Result<ExecutionPlan, EngineError> {
920 if options.max_depth == 0 {
921 return Err(EngineError::InvalidWorkflow(
922 "max_depth must be at least 1".to_string(),
923 ));
924 }
925
926 let handler = self
927 .handlers
928 .get(handler_name)
929 .ok_or_else(|| {
930 EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
931 })?
932 .clone();
933
934 let estimates = if options.estimate_durations {
935 estimate_durations(&self.store, handler_name, options.sample_runs).await?
936 } else {
937 HashMap::new()
938 };
939
940 let shared: SharedPlanRecorder = Arc::new(Mutex::new(PlanRecorder::new(
941 handler_name.to_string(),
942 payload,
943 options.max_depth,
944 estimates,
945 )));
946
947 let handlers = self.handlers.clone();
950 let resolver: crate::context::HandlerResolver =
951 Arc::new(move |name: &str| handlers.get(name).cloned());
952 let mut ctx = WorkflowContext::with_handler_resolver(
953 Uuid::now_v7(),
954 handler_name.to_string(),
955 self.store.clone(),
956 self.provider.clone(),
957 resolver,
958 );
959 ctx.set_plan(shared.clone());
960
961 if let Err(err) = handler.execute(&mut ctx).await {
962 lock_plan(&shared).fail(err.to_string());
963 }
964 drop(ctx);
965
966 let plan = match Arc::try_unwrap(shared) {
967 Ok(mutex) => mutex
968 .into_inner()
969 .unwrap_or_else(|poisoned| poisoned.into_inner())
970 .into_plan(),
971 Err(shared) => lock_plan(&shared).snapshot(),
972 };
973
974 info!(
975 workflow = %handler_name,
976 steps = plan.steps.len(),
977 truncated = plan.truncated,
978 "execution plan built"
979 );
980
981 Ok(plan)
982 }
983
984 #[tracing::instrument(name = "engine.enqueue_handler", skip_all, fields(workflow = %handler_name))]
995 pub async fn enqueue_handler(
996 &self,
997 handler_name: &str,
998 trigger: TriggerKind,
999 payload: Value,
1000 max_retries: u32,
1001 ) -> Result<Run, EngineError> {
1002 self.enqueue_handler_with_options(
1003 handler_name,
1004 trigger,
1005 payload,
1006 EnqueueOptions {
1007 max_retries,
1008 ..Default::default()
1009 },
1010 )
1011 .await
1012 .map(RunCreation::into_run)
1013 }
1014
1015 #[tracing::instrument(name = "engine.enqueue_handler_with_options", skip_all, fields(workflow = %handler_name))]
1060 pub async fn enqueue_handler_with_options(
1061 &self,
1062 handler_name: &str,
1063 trigger: TriggerKind,
1064 payload: Value,
1065 options: EnqueueOptions,
1066 ) -> Result<RunCreation, EngineError> {
1067 let EnqueueOptions {
1068 max_retries,
1069 labels,
1070 scheduled_at,
1071 max_cost_usd,
1072 created_by,
1073 idempotency_key,
1074 } = options;
1075
1076 let handler = self.handlers.get(handler_name).ok_or_else(|| {
1077 EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
1078 })?;
1079
1080 self.check_monthly_quota(handler_name).await?;
1081
1082 let handler_version = handler.version().map(str::to_string);
1083 let mut merged_labels = handler.default_labels();
1084 merged_labels.extend(labels);
1085 let resolved_cap = self
1086 .budget
1087 .resolve_run_cap(max_cost_usd, handler.default_max_cost_usd());
1088
1089 let creation = self
1090 .store
1091 .create_run(NewRun {
1092 workflow_name: handler_name.to_string(),
1093 trigger,
1094 payload,
1095 max_retries,
1096 handler_version,
1097 labels: merged_labels,
1098 scheduled_at,
1099 created_by,
1100 idempotency_key,
1101 max_cost_usd: resolved_cap,
1102 })
1103 .await?;
1104
1105 match &creation {
1106 RunCreation::Created(run) => info!(
1107 run_id = %run.id,
1108 workflow = %handler_name,
1109 max_cost_usd = ?resolved_cap,
1110 "handler run enqueued"
1111 ),
1112 RunCreation::Existing(run) => info!(
1113 run_id = %run.id,
1114 workflow = %handler_name,
1115 "idempotent replay, nothing enqueued"
1116 ),
1117 }
1118
1119 Ok(creation)
1120 }
1121
1122 #[tracing::instrument(name = "engine.execute_handler_run", skip_all, fields(run_id = %run_id))]
1142 pub async fn execute_handler_run(&self, run_id: Uuid) -> Result<WorkflowResult, EngineError> {
1143 let run = self
1144 .store
1145 .get_run(run_id)
1146 .await?
1147 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
1148
1149 if let Some(root_run_id) = chain_root(&run) {
1150 return self.resume_chain(run_id, root_run_id).await;
1151 }
1152
1153 let handler = self
1154 .handlers
1155 .get(&run.workflow_name)
1156 .ok_or_else(|| {
1157 EngineError::InvalidWorkflow(format!(
1158 "no handler registered: {}",
1159 run.workflow_name
1160 ))
1161 })?
1162 .clone();
1163
1164 #[cfg(feature = "prometheus")]
1165 gauge!(RUNS_ACTIVE, "workflow" => run.workflow_name.clone()).increment(1.0);
1166
1167 let run_start = Instant::now();
1168 let mut ctx = self.build_context_with_guard(&run, handler.as_ref());
1169
1170 let result = if handler.is_version_compatible(run.handler_version.as_deref()) {
1181 ctx.load_replay_steps().await?;
1182 self.release_then_execute(run_id, handler.as_ref(), &mut ctx)
1183 .await
1184 } else {
1185 Err(EngineError::HandlerVersionMismatch {
1186 run_id,
1187 workflow_name: run.workflow_name.clone(),
1188 run_version: run
1189 .handler_version
1190 .clone()
1191 .unwrap_or_else(|| "unknown".to_string()),
1192 current_version: handler
1193 .version()
1194 .map(str::to_string)
1195 .unwrap_or_else(|| "unknown".to_string()),
1196 })
1197 };
1198
1199 self.finalize_run(
1200 run_id,
1201 &run.workflow_name,
1202 result,
1203 &ctx,
1204 run_start,
1205 run.labels,
1206 )
1207 .await
1208 }
1209
1210 #[tracing::instrument(name = "engine.execute_run", skip_all, fields(run_id = %run_id))]
1218 pub async fn execute_run(&self, run_id: Uuid) -> Result<WorkflowResult, EngineError> {
1219 self.execute_handler_run(run_id).await
1220 }
1221
1222 #[tracing::instrument(name = "engine.resume_run", skip_all, fields(run_id = %run_id))]
1251 pub async fn resume_run(&self, run_id: Uuid) -> Result<WorkflowResult, EngineError> {
1252 let run = self
1253 .store
1254 .get_run(run_id)
1255 .await?
1256 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
1257
1258 if let Some(root_run_id) = chain_root(&run) {
1259 return self.resume_chain(run_id, root_run_id).await;
1260 }
1261
1262 self.resume_loaded_run(run).await
1263 }
1264
1265 async fn resume_chain(
1273 &self,
1274 child_run_id: Uuid,
1275 root_run_id: Uuid,
1276 ) -> Result<WorkflowResult, EngineError> {
1277 let root = self
1278 .store
1279 .get_run(root_run_id)
1280 .await?
1281 .ok_or(EngineError::Store(StoreError::RunNotFound(root_run_id)))?;
1282
1283 match root.status.state {
1284 RunStatus::AwaitingApproval | RunStatus::Pending => {
1285 self.store
1286 .update_run_status(root_run_id, RunStatus::Running)
1287 .await?;
1288 }
1289 RunStatus::Sleeping => {
1290 self.store
1291 .update_run_status(root_run_id, RunStatus::Pending)
1292 .await?;
1293 self.store
1294 .update_run_status(root_run_id, RunStatus::Running)
1295 .await?;
1296 }
1297 other => {
1298 let reason = format!(
1299 "cannot resume child run {child_run_id}: root run {root_run_id} is {other}"
1300 );
1301 if let Err(err) = self
1302 .fail_or_schedule_retry(child_run_id, &reason, false, None, None)
1303 .await
1304 {
1305 error!(
1306 run_id = %child_run_id,
1307 error = %err,
1308 "failed to fail a child run whose root cannot resume"
1309 );
1310 }
1311 return Err(EngineError::InvalidWorkflow(reason));
1312 }
1313 }
1314
1315 info!(
1316 run_id = %child_run_id,
1317 root_run_id = %root_run_id,
1318 "child run resumed through its root run"
1319 );
1320
1321 let root = self
1322 .store
1323 .get_run(root_run_id)
1324 .await?
1325 .ok_or(EngineError::Store(StoreError::RunNotFound(root_run_id)))?;
1326 self.resume_loaded_run(root).await
1327 }
1328
1329 async fn resume_loaded_run(&self, run: Run) -> Result<WorkflowResult, EngineError> {
1331 let run_id = run.id;
1332 let handler = self
1333 .handlers
1334 .get(&run.workflow_name)
1335 .ok_or_else(|| {
1336 EngineError::InvalidWorkflow(format!(
1337 "no handler registered: {}",
1338 run.workflow_name
1339 ))
1340 })?
1341 .clone();
1342
1343 info!(run_id = %run_id, workflow = %run.workflow_name, "resuming run after approval");
1344
1345 let run_start = Instant::now();
1346 let mut ctx = self.build_context_with_guard(&run, handler.as_ref());
1347
1348 let result = if handler.is_version_compatible(run.handler_version.as_deref()) {
1349 ctx.load_replay_steps().await?;
1350 self.release_then_execute(run_id, handler.as_ref(), &mut ctx)
1351 .await
1352 } else {
1353 Err(EngineError::HandlerVersionMismatch {
1354 run_id,
1355 workflow_name: run.workflow_name.clone(),
1356 run_version: run
1357 .handler_version
1358 .clone()
1359 .unwrap_or_else(|| "unknown".to_string()),
1360 current_version: handler
1361 .version()
1362 .map(str::to_string)
1363 .unwrap_or_else(|| "unknown".to_string()),
1364 })
1365 };
1366
1367 self.finalize_run(
1368 run_id,
1369 &run.workflow_name,
1370 result,
1371 &ctx,
1372 run_start,
1373 run.labels,
1374 )
1375 .await
1376 }
1377
1378 pub async fn deliver_signal(
1424 self: &Arc<Self>,
1425 signal: NewSignal,
1426 ) -> Result<SignalDelivery, EngineError> {
1427 if signal.name.trim().is_empty() {
1428 return Err(EngineError::InvalidSignal(
1429 "signal name must not be empty".to_string(),
1430 ));
1431 }
1432 if signal.key.trim().is_empty() {
1433 return Err(EngineError::InvalidSignal(
1434 "signal key must not be empty".to_string(),
1435 ));
1436 }
1437
1438 let stored = match self.store.insert_signal(signal).await? {
1439 SignalInsert::Created(stored) => stored,
1440 SignalInsert::Duplicate(existing) => {
1441 info!(
1442 signal_id = %existing.id,
1443 signal = %existing.name,
1444 key = %existing.key,
1445 "duplicate signal ignored"
1446 );
1447 return Ok(SignalDelivery {
1448 signal_id: existing.id,
1449 duplicate: true,
1450 resumed: Vec::new(),
1451 rejected: Vec::new(),
1452 });
1453 }
1454 };
1455
1456 let waiters = self
1457 .store
1458 .list_signal_waiters(&stored.name, &stored.key)
1459 .await?;
1460 let mut resumed = Vec::new();
1461 let mut rejected = Vec::new();
1462
1463 for step in waiters {
1464 if let Err(error) = validate_step_payload(step.input.as_ref(), &stored.payload) {
1465 rejected.push(SignalRejected {
1466 run_id: step.run_id,
1467 step_id: step.id,
1468 error,
1469 });
1470 continue;
1471 }
1472
1473 match self
1474 .store
1475 .resolve_signal_step(step.id, received_output(&stored))
1476 .await
1477 {
1478 Ok(SignalStepResolution::Resolved {
1479 run_id,
1480 run_resumed,
1481 }) => {
1482 resumed.push(SignalResumed {
1483 run_id,
1484 step_id: step.id,
1485 });
1486 if run_resumed && self.execution_mode == ExecutionMode::Local {
1487 self.spawn_local_resume(run_id);
1488 }
1489 }
1490 Ok(SignalStepResolution::NotWaiting { .. }) => {}
1492 Err(err) => {
1493 error!(
1494 run_id = %step.run_id,
1495 step_id = %step.id,
1496 error = %err,
1497 "failed to resolve a waiting signal step"
1498 );
1499 rejected.push(SignalRejected {
1500 run_id: step.run_id,
1501 step_id: step.id,
1502 error: err.to_string(),
1503 });
1504 }
1505 }
1506 }
1507
1508 info!(
1509 signal_id = %stored.id,
1510 signal = %stored.name,
1511 key = %stored.key,
1512 resumed = resumed.len(),
1513 rejected = rejected.len(),
1514 "signal received"
1515 );
1516 self.event_publisher
1517 .publish(Event::SignalReceived(SignalReceivedEvent {
1518 signal_id: stored.id,
1519 name: stored.name.clone(),
1520 key: stored.key.clone(),
1521 resumed_runs: resumed.iter().map(|r| r.run_id).collect(),
1522 at: stored.received_at,
1523 }));
1524
1525 Ok(SignalDelivery {
1526 signal_id: stored.id,
1527 duplicate: false,
1528 resumed,
1529 rejected,
1530 })
1531 }
1532
1533 pub async fn send_signal<S: Signal>(
1571 self: &Arc<Self>,
1572 signal: &S,
1573 key: &str,
1574 idempotency_id: Option<&str>,
1575 ) -> Result<SignalDelivery, EngineError> {
1576 let payload = to_value(signal)?;
1577 self.deliver_signal(NewSignal {
1578 name: S::NAME.to_string(),
1579 key: key.to_string(),
1580 payload,
1581 idempotency_id: idempotency_id.map(str::to_string),
1582 })
1583 .await
1584 }
1585
1586 pub(crate) fn spawn_local_resume(self: &Arc<Self>, run_id: Uuid) {
1592 let engine = Arc::clone(self);
1593 spawn(async move {
1594 if let Err(err) = engine
1595 .store
1596 .update_run_status(run_id, RunStatus::Running)
1597 .await
1598 {
1599 error!(run_id = %run_id, error = %err, "failed to restart a woken run");
1600 return;
1601 }
1602 if let Err(err) = engine.resume_run(run_id).await {
1603 error!(run_id = %run_id, error = %err, "failed to resume a woken run");
1604 }
1605 });
1606 }
1607
1608 pub async fn fail_or_schedule_retry(
1653 &self,
1654 run_id: Uuid,
1655 error: &str,
1656 retryable: bool,
1657 cost_usd: Option<Decimal>,
1658 duration_ms: Option<u64>,
1659 ) -> Result<RunStatus, EngineError> {
1660 let run = self
1661 .store
1662 .get_run(run_id)
1663 .await?
1664 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
1665
1666 let has_attempts_left = run.retry_count < run.max_retries;
1667 let update = if retryable && has_attempts_left {
1668 let backoff = backoff_for_retry(run.retry_count);
1669 let scheduled_at = Utc::now() + TimeDelta::milliseconds(backoff.as_millis() as i64);
1670
1671 info!(
1672 run_id = %run_id,
1673 workflow = %run.workflow_name,
1674 attempt = run.retry_count + 1,
1675 max_retries = run.max_retries,
1676 backoff_secs = backoff.as_secs(),
1677 scheduled_at = %scheduled_at,
1678 "run failed, scheduling retry"
1679 );
1680
1681 RunUpdate {
1682 status: Some(RunStatus::Retrying),
1683 error: Some(error.to_string()),
1684 increment_retry: true,
1685 cost_usd,
1686 duration_ms,
1687 scheduled_at: Some(scheduled_at),
1688 ..RunUpdate::default()
1689 }
1690 } else {
1691 RunUpdate {
1692 status: Some(RunStatus::Failed),
1693 error: Some(error.to_string()),
1694 cost_usd,
1695 duration_ms,
1696 completed_at: Some(Utc::now()),
1697 ..RunUpdate::default()
1698 }
1699 };
1700
1701 let status = update.status.unwrap_or(RunStatus::Failed);
1702 self.store.update_run(run_id, update).await?;
1703 self.fail_orphaned_steps(run_id, error).await?;
1704
1705 Ok(status)
1706 }
1707
1708 pub async fn fail_orphaned_steps(
1722 &self,
1723 run_id: Uuid,
1724 error_message: &str,
1725 ) -> Result<(), EngineError> {
1726 let steps = self.store.list_steps(run_id).await?;
1727 let now = Utc::now();
1728
1729 for step in steps {
1730 if step.status.state.is_terminal() {
1731 continue;
1732 }
1733
1734 let (target_status, error) = match step.status.state {
1735 StepStatus::Running | StepStatus::AwaitingApproval => {
1736 let err = if step.error.is_some() {
1737 None
1738 } else {
1739 Some(error_message.to_string())
1740 };
1741 (StepStatus::Failed, err)
1742 }
1743 StepStatus::Pending => (StepStatus::Skipped, None),
1744 _ => continue,
1745 };
1746
1747 if let Err(e) = self
1748 .store
1749 .update_step(
1750 step.id,
1751 StepUpdate {
1752 status: Some(target_status),
1753 error,
1754 completed_at: Some(now),
1755 ..StepUpdate::default()
1756 },
1757 )
1758 .await
1759 {
1760 warn!(
1761 run_id = %run_id,
1762 step_id = %step.id,
1763 step_name = %step.name,
1764 error = %e,
1765 "failed to cleanup orphaned step"
1766 );
1767 } else {
1768 info!(
1769 run_id = %run_id,
1770 step_id = %step.id,
1771 step_name = %step.name,
1772 from = %step.status.state,
1773 to = %target_status,
1774 "cleaned up orphaned step"
1775 );
1776 }
1777 }
1778
1779 Ok(())
1780 }
1781
1782 async fn release_then_execute(
1787 &self,
1788 run_id: Uuid,
1789 handler: &dyn WorkflowHandler,
1790 ctx: &mut WorkflowContext,
1791 ) -> Result<(), EngineError> {
1792 match self.provider.release_run(&run_id.to_string()).await {
1793 Ok(()) => handler.execute(ctx).await,
1794 Err(e) => Err(EngineError::Operation(OperationError::Agent(e))),
1795 }
1796 }
1797
1798 async fn finalize_run(
1804 &self,
1805 run_id: Uuid,
1806 workflow_name: &str,
1807 result: Result<(), EngineError>,
1808 ctx: &WorkflowContext,
1809 run_start: Instant,
1810 run_labels: HashMap<String, String>,
1811 ) -> Result<WorkflowResult, EngineError> {
1812 let total_duration = ctx.carried_duration_ms() + run_start.elapsed().as_millis() as u64;
1815 let completed_at = Utc::now();
1816
1817 let final_status;
1818 let final_run;
1819
1820 match result {
1821 Ok(()) => {
1822 final_status = if ctx.has_allowed_failure() {
1823 RunStatus::Warning
1824 } else {
1825 RunStatus::Completed
1826 };
1827 final_run = self
1828 .store
1829 .update_run_returning(
1830 run_id,
1831 RunUpdate {
1832 status: Some(final_status),
1833 cost_usd: Some(ctx.total_cost_usd()),
1834 duration_ms: Some(total_duration),
1835 completed_at: Some(completed_at),
1836 ..RunUpdate::default()
1837 },
1838 )
1839 .await?;
1840
1841 info!(
1842 run_id = %run_id,
1843 status = %final_status,
1844 cost_usd = %ctx.total_cost_usd(),
1845 duration_ms = total_duration,
1846 "run completed"
1847 );
1848 }
1849 Err(EngineError::ApprovalRequired {
1850 run_id: approval_run_id,
1851 step_id,
1852 ref message,
1853 }) => {
1854 final_status = RunStatus::AwaitingApproval;
1855 final_run = self
1856 .store
1857 .update_run_returning(
1858 run_id,
1859 RunUpdate {
1860 status: Some(RunStatus::AwaitingApproval),
1861 cost_usd: Some(ctx.total_cost_usd()),
1862 duration_ms: Some(total_duration),
1863 ..RunUpdate::default()
1864 },
1865 )
1866 .await?;
1867
1868 info!(
1869 run_id = %approval_run_id,
1870 step_id = %step_id,
1871 message = %message,
1872 "run awaiting approval"
1873 );
1874
1875 self.publish_approval_requested(approval_run_id, step_id, message)
1876 .await?;
1877 }
1878 Err(EngineError::ChildSuspended {
1879 run_id: child_run_id,
1880 ref cause,
1881 }) => {
1882 final_status = cause.suspension_status();
1883 final_run = self
1887 .store
1888 .update_run_returning(
1889 run_id,
1890 RunUpdate {
1891 status: Some(final_status),
1892 cost_usd: Some(ctx.total_cost_usd()),
1893 duration_ms: Some(total_duration),
1894 ..RunUpdate::default()
1895 },
1896 )
1897 .await?;
1898
1899 let leaf = cause.suspension_leaf();
1900 info!(
1901 run_id = %run_id,
1902 child_run_id = %child_run_id,
1903 status = %final_status,
1904 cause = %leaf,
1905 "run suspended with its child run"
1906 );
1907
1908 match leaf {
1909 EngineError::ApprovalRequired {
1910 run_id: approval_run_id,
1911 step_id,
1912 message,
1913 } => {
1914 self.publish_approval_requested(*approval_run_id, *step_id, message)
1915 .await?;
1916 }
1917 EngineError::SignalWaiting {
1918 run_id: wait_run_id,
1919 step_id,
1920 step_name,
1921 name,
1922 key,
1923 deadline_at,
1924 } => {
1925 self.event_publisher
1926 .publish(Event::SignalAwaited(SignalAwaitedEvent {
1927 run_id: *wait_run_id,
1928 step_id: *step_id,
1929 step_name: step_name.clone(),
1930 name: name.clone(),
1931 key: key.clone(),
1932 deadline_at: *deadline_at,
1933 at: Utc::now(),
1934 }));
1935 }
1936 _ => {}
1939 }
1940 }
1941 Err(EngineError::HumanInputRequired {
1942 run_id: input_run_id,
1943 step_id,
1944 ref message,
1945 }) => {
1946 final_status = RunStatus::AwaitingApproval;
1947 final_run = self
1948 .store
1949 .update_run_returning(
1950 run_id,
1951 RunUpdate {
1952 status: Some(RunStatus::AwaitingApproval),
1953 cost_usd: Some(ctx.total_cost_usd()),
1954 duration_ms: Some(total_duration),
1955 ..RunUpdate::default()
1956 },
1957 )
1958 .await?;
1959
1960 info!(
1962 run_id = %input_run_id,
1963 step_id = %step_id,
1964 message = %message,
1965 "run awaiting human input"
1966 );
1967 }
1968 Err(EngineError::DelaySleeping {
1969 run_id: delay_run_id,
1970 step_id,
1971 wake_at,
1972 }) => {
1973 final_status = RunStatus::Sleeping;
1974 final_run = self
1975 .store
1976 .update_run_returning(
1977 run_id,
1978 RunUpdate {
1979 status: Some(RunStatus::Sleeping),
1980 cost_usd: Some(ctx.total_cost_usd()),
1981 duration_ms: Some(total_duration),
1982 scheduled_at: Some(wake_at),
1983 ..RunUpdate::default()
1984 },
1985 )
1986 .await?;
1987
1988 info!(
1989 run_id = %delay_run_id,
1990 step_id = %step_id,
1991 wake_at = %wake_at,
1992 "run sleeping until delay elapses"
1993 );
1994 }
1995 Err(EngineError::SignalWaiting {
1996 run_id: wait_run_id,
1997 step_id,
1998 ref step_name,
1999 ref name,
2000 ref key,
2001 deadline_at,
2002 }) => {
2003 final_status = RunStatus::Sleeping;
2004 let waiting = self
2008 .store
2009 .suspend_run_on_signal(run_id, step_id, deadline_at)
2010 .await?;
2011 final_run = self
2012 .store
2013 .update_run_returning(
2014 run_id,
2015 RunUpdate {
2016 cost_usd: Some(ctx.total_cost_usd()),
2017 duration_ms: Some(total_duration),
2018 ..RunUpdate::default()
2019 },
2020 )
2021 .await?;
2022
2023 if waiting {
2024 self.event_publisher
2025 .publish(Event::SignalAwaited(SignalAwaitedEvent {
2026 run_id: wait_run_id,
2027 step_id,
2028 step_name: step_name.clone(),
2029 name: name.clone(),
2030 key: key.clone(),
2031 deadline_at,
2032 at: Utc::now(),
2033 }));
2034 }
2035
2036 info!(
2037 run_id = %wait_run_id,
2038 step_id = %step_id,
2039 signal = %name,
2040 key = %key,
2041 deadline_at = %deadline_at,
2042 waiting,
2043 "run sleeping until a signal arrives"
2044 );
2045 }
2046 Err(err) => {
2047 let guardrail_stop = matches!(
2051 err,
2052 EngineError::RunBudgetExceeded { .. } | EngineError::WorkflowGuardRejected(_)
2053 );
2054
2055 final_status = if guardrail_stop {
2056 if let Err(store_err) = self
2057 .store
2058 .update_run(
2059 run_id,
2060 RunUpdate {
2061 status: Some(RunStatus::Cancelled),
2062 error: Some(err.to_string()),
2063 cost_usd: Some(ctx.total_cost_usd()),
2064 duration_ms: Some(total_duration),
2065 completed_at: Some(completed_at),
2066 ..RunUpdate::default()
2067 },
2068 )
2069 .await
2070 {
2071 error!(run_id = %run_id, store_error = %store_err, "failed to persist run cancellation");
2072 }
2073 if let Err(cleanup_err) = self
2074 .fail_orphaned_steps(run_id, "run stopped: guardrail limit reached")
2075 .await
2076 {
2077 error!(run_id = %run_id, store_error = %cleanup_err, "failed to cleanup orphaned steps");
2078 }
2079 RunStatus::Cancelled
2080 } else {
2081 self.fail_or_schedule_retry(
2082 run_id,
2083 &err.to_string(),
2084 is_run_retryable(&err),
2085 Some(ctx.total_cost_usd()),
2086 Some(total_duration),
2087 )
2088 .await
2089 .unwrap_or_else(|store_err| {
2090 error!(run_id = %run_id, store_error = %store_err, "failed to persist run failure");
2091 RunStatus::Failed
2092 })
2093 };
2094
2095 if matches!(err, EngineError::RunBudgetExceeded { .. }) {
2096 self.on_run_budget_exceeded(workflow_name, run_id, &err);
2097 }
2098
2099 error!(run_id = %run_id, status = %final_status, error = %err, "run stopped");
2100
2101 self.publish_run_status_changed(
2102 workflow_name,
2103 run_id,
2104 final_status,
2105 Some(err.to_string()),
2106 ctx,
2107 total_duration,
2108 run_labels,
2109 );
2110
2111 #[cfg(feature = "prometheus")]
2112 self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
2113
2114 return Err(err);
2115 }
2116 }
2117
2118 self.publish_run_status_changed(
2119 workflow_name,
2120 run_id,
2121 final_status,
2122 None,
2123 ctx,
2124 total_duration,
2125 run_labels,
2126 );
2127
2128 #[cfg(feature = "prometheus")]
2129 self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
2130
2131 Ok(WorkflowResult {
2132 run: final_run,
2133 steps: ctx.step_results().to_vec(),
2134 })
2135 }
2136
2137 async fn publish_approval_requested(
2140 &self,
2141 run_id: Uuid,
2142 step_id: Uuid,
2143 message: &str,
2144 ) -> Result<(), EngineError> {
2145 let requirement = self
2146 .store
2147 .get_step(step_id)
2148 .await?
2149 .and_then(|s| s.approval_requirement);
2150 self.event_publisher
2151 .publish(Event::ApprovalRequested(ApprovalRequestedEvent {
2152 run_id,
2153 step_id,
2154 message: message.to_string(),
2155 requirement,
2156 at: Utc::now(),
2157 }));
2158 Ok(())
2159 }
2160
2161 pub async fn fail_ancestors(&self, run_id: Uuid, reason: &str) -> Result<(), EngineError> {
2190 let mut current = self
2191 .store
2192 .get_run(run_id)
2193 .await?
2194 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
2195 let mut visited = HashSet::from([run_id]);
2197
2198 while let Some(parent_id) = chain_parent(¤t) {
2199 if !visited.insert(parent_id) {
2200 break;
2201 }
2202 let status = self
2203 .fail_or_schedule_retry(parent_id, reason, false, None, None)
2204 .await?;
2205 info!(
2206 run_id = %run_id,
2207 ancestor_run_id = %parent_id,
2208 status = %status,
2209 "ancestor run failed with its child"
2210 );
2211 current = self
2212 .store
2213 .get_run(parent_id)
2214 .await?
2215 .ok_or(EngineError::Store(StoreError::RunNotFound(parent_id)))?;
2216 }
2217
2218 Ok(())
2219 }
2220
2221 #[cfg(feature = "prometheus")]
2223 fn emit_run_metrics(
2224 &self,
2225 workflow_name: &str,
2226 status: RunStatus,
2227 duration_ms: u64,
2228 ctx: &WorkflowContext,
2229 ) {
2230 let status_str = status.to_string();
2231 let wf = workflow_name.to_string();
2232
2233 counter!(RUNS_TOTAL, "workflow" => wf.clone(), "status" => status_str.clone()).increment(1);
2234 histogram!(RUN_DURATION_SECONDS, "workflow" => wf.clone(), "status" => status_str)
2235 .record(duration_ms as f64 / 1000.0);
2236 histogram!(RUN_COST_USD, "workflow" => wf.clone()).record(
2237 ctx.total_cost_usd()
2238 .to_string()
2239 .parse::<f64>()
2240 .unwrap_or(0.0),
2241 );
2242 gauge!(RUNS_ACTIVE, "workflow" => wf).decrement(1.0);
2243 }
2244
2245 fn on_run_budget_exceeded(&self, workflow_name: &str, run_id: Uuid, err: &EngineError) {
2251 let EngineError::RunBudgetExceeded {
2252 limit_usd,
2253 spent_usd,
2254 step_budget_usd,
2255 ..
2256 } = err
2257 else {
2258 return;
2259 };
2260
2261 #[cfg(feature = "prometheus")]
2262 counter!(
2263 RUN_BUDGET_EXCEEDED_TOTAL,
2264 "workflow" => workflow_name.to_string(),
2265 "scope" => "run",
2266 )
2267 .increment(1);
2268
2269 self.event_publisher
2270 .publish(Event::RunBudgetExceeded(RunBudgetExceededEvent {
2271 run_id,
2272 workflow_name: workflow_name.to_string(),
2273 limit_usd: *limit_usd,
2274 spent_usd: *spent_usd,
2275 step_budget_usd: *step_budget_usd,
2276 at: Utc::now(),
2277 }));
2278 }
2279
2280 #[allow(clippy::too_many_arguments)]
2285 fn publish_run_status_changed(
2286 &self,
2287 workflow_name: &str,
2288 run_id: Uuid,
2289 to: RunStatus,
2290 error: Option<String>,
2291 ctx: &WorkflowContext,
2292 duration_ms: u64,
2293 labels: HashMap<String, String>,
2294 ) {
2295 let now = Utc::now();
2296 let cost_usd = ctx.total_cost_usd();
2297 let wf = workflow_name.to_string();
2298
2299 self.event_publisher
2300 .publish(Event::RunStatusChanged(RunStatusChangedEvent {
2301 run_id,
2302 workflow_name: wf.clone(),
2303 from: RunStatus::Running,
2304 to,
2305 error: error.clone(),
2306 cost_usd,
2307 duration_ms,
2308 labels: labels.clone(),
2309 at: now,
2310 }));
2311
2312 if to == RunStatus::Failed {
2313 self.event_publisher
2314 .publish(Event::RunFailed(RunFailedEvent {
2315 run_id,
2316 workflow_name: wf,
2317 error,
2318 cost_usd,
2319 duration_ms,
2320 labels,
2321 at: now,
2322 }));
2323 }
2324 }
2325}
2326
2327impl fmt::Debug for Engine {
2328 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
2329 f.debug_struct("Engine")
2330 .field("handlers", &self.handlers.keys().collect::<Vec<_>>())
2331 .finish_non_exhaustive()
2332 }
2333}
2334
2335#[cfg(test)]
2336mod tests {
2337 use super::*;
2338 use crate::config::ShellConfig;
2339 use crate::handler::{HandlerFuture, WorkflowHandler};
2340 use ironflow_core::providers::claude::ClaudeCodeProvider;
2341 use ironflow_core::providers::record_replay::RecordReplayProvider;
2342 use ironflow_store::memory::InMemoryStore;
2343 use ironflow_store::models::StepStatus;
2344 use serde_json::json;
2345
2346 struct EchoWorkflow;
2348
2349 impl WorkflowHandler for EchoWorkflow {
2350 fn name(&self) -> &str {
2351 "echo-workflow"
2352 }
2353
2354 fn describe(&self) -> WorkflowInfo {
2355 WorkflowInfo {
2356 description: "A simple workflow that echoes hello".to_string(),
2357 source_code: None,
2358 sub_workflows: Vec::new(),
2359 category: None,
2360 version: self.version().map(str::to_string),
2361 compatible_versions: Vec::new(),
2362 input_schema: None,
2363 default_labels: HashMap::new(),
2364 schedule: self.schedule().cloned(),
2365 default_max_cost_usd: self.default_max_cost_usd(),
2366 }
2367 }
2368
2369 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2370 Box::pin(async move {
2371 ctx.shell("greet", ShellConfig::new("echo hello")).await?;
2372 Ok(())
2373 })
2374 }
2375 }
2376
2377 struct FailingWorkflow;
2379
2380 impl WorkflowHandler for FailingWorkflow {
2381 fn name(&self) -> &str {
2382 "failing-workflow"
2383 }
2384
2385 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2386 Box::pin(async move {
2387 ctx.shell("fail", ShellConfig::new("exit 1")).await?;
2388 Ok(())
2389 })
2390 }
2391 }
2392
2393 fn create_test_engine() -> Engine {
2394 let store = Arc::new(InMemoryStore::new());
2395 let inner = ClaudeCodeProvider::new();
2396 let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
2397 inner,
2398 "/tmp/ironflow-fixtures",
2399 ));
2400 Engine::new(store, provider)
2401 }
2402
2403 #[test]
2404 fn engine_new_creates_instance() {
2405 let engine = create_test_engine();
2406 assert_eq!(engine.handler_names().len(), 0);
2407 }
2408
2409 #[test]
2410 fn execution_mode_defaults_to_local() {
2411 let engine = create_test_engine();
2412 assert_eq!(engine.execution_mode(), ExecutionMode::Local);
2413 }
2414
2415 #[test]
2416 fn with_execution_mode_overrides_the_default() {
2417 let engine = create_test_engine().with_execution_mode(ExecutionMode::Workers);
2418 assert_eq!(engine.execution_mode(), ExecutionMode::Workers);
2419 }
2420
2421 #[test]
2422 fn engine_register_handler() {
2423 let mut engine = create_test_engine();
2424 let result = engine.register(EchoWorkflow);
2425 assert!(result.is_ok());
2426 assert_eq!(engine.handler_names().len(), 1);
2427 assert!(engine.handler_names().contains(&"echo-workflow"));
2428 }
2429
2430 #[test]
2431 fn engine_register_duplicate_returns_error() {
2432 let mut engine = create_test_engine();
2433 engine.register(EchoWorkflow).unwrap();
2434 let result = engine.register(EchoWorkflow);
2435 assert!(result.is_err());
2436 }
2437
2438 #[test]
2439 fn engine_get_handler_found() {
2440 let mut engine = create_test_engine();
2441 engine.register(EchoWorkflow).unwrap();
2442 let handler = engine.get_handler("echo-workflow");
2443 assert!(handler.is_some());
2444 }
2445
2446 #[test]
2447 fn engine_get_handler_not_found() {
2448 let engine = create_test_engine();
2449 let handler = engine.get_handler("nonexistent");
2450 assert!(handler.is_none());
2451 }
2452
2453 #[test]
2454 fn engine_handler_names_lists_all() {
2455 let mut engine = create_test_engine();
2456 engine.register(EchoWorkflow).unwrap();
2457 engine.register(FailingWorkflow).unwrap();
2458 let names = engine.handler_names();
2459 assert_eq!(names.len(), 2);
2460 assert!(names.contains(&"echo-workflow"));
2461 assert!(names.contains(&"failing-workflow"));
2462 }
2463
2464 #[test]
2465 fn engine_handler_info_returns_description() {
2466 let mut engine = create_test_engine();
2467 engine.register(EchoWorkflow).unwrap();
2468 let info = engine.handler_info("echo-workflow");
2469 assert!(info.is_some());
2470 let info = info.unwrap();
2471 assert_eq!(info.description, "A simple workflow that echoes hello");
2472 }
2473
2474 struct CategorizedWorkflow;
2475
2476 impl WorkflowHandler for CategorizedWorkflow {
2477 fn name(&self) -> &str {
2478 "categorized"
2479 }
2480 fn category(&self) -> Option<&str> {
2481 Some("data/etl")
2482 }
2483 fn execute<'a>(
2484 &'a self,
2485 _ctx: &'a mut WorkflowContext,
2486 ) -> crate::handler::HandlerFuture<'a> {
2487 Box::pin(async move { Ok(()) })
2488 }
2489 }
2490
2491 #[test]
2492 fn engine_default_describe_propagates_category() {
2493 let mut engine = create_test_engine();
2494 engine.register(CategorizedWorkflow).unwrap();
2495 let info = engine.handler_info("categorized").unwrap();
2496 assert_eq!(info.category.as_deref(), Some("data/etl"));
2497 }
2498
2499 #[test]
2500 fn engine_default_describe_without_category() {
2501 let mut engine = create_test_engine();
2502 engine.register(EchoWorkflow).unwrap();
2503 let info = engine.handler_info("echo-workflow").unwrap();
2504 assert!(info.category.is_none());
2505 }
2506
2507 struct ScheduledWorkflow {
2512 schedule: CronSchedule,
2513 }
2514
2515 impl ScheduledWorkflow {
2516 fn new() -> Self {
2517 Self {
2518 schedule: CronSchedule::new("0 0 * * * *").unwrap(),
2519 }
2520 }
2521 }
2522
2523 impl WorkflowHandler for ScheduledWorkflow {
2524 fn name(&self) -> &str {
2525 "scheduled"
2526 }
2527 fn schedule(&self) -> Option<&CronSchedule> {
2528 Some(&self.schedule)
2529 }
2530 fn execute<'a>(
2531 &'a self,
2532 _ctx: &'a mut WorkflowContext,
2533 ) -> crate::handler::HandlerFuture<'a> {
2534 Box::pin(async move { Ok(()) })
2535 }
2536 }
2537
2538 #[test]
2539 fn engine_default_describe_propagates_schedule() {
2540 let mut engine = create_test_engine();
2541 engine.register(ScheduledWorkflow::new()).unwrap();
2542 let info = engine.handler_info("scheduled").unwrap();
2543 assert_eq!(
2544 info.schedule.as_ref().map(|s| s.as_str()),
2545 Some("0 0 * * * *")
2546 );
2547 }
2548
2549 #[test]
2550 fn engine_default_describe_without_schedule() {
2551 let mut engine = create_test_engine();
2552 engine.register(EchoWorkflow).unwrap();
2553 let info = engine.handler_info("echo-workflow").unwrap();
2554 assert!(info.schedule.is_none());
2555 }
2556
2557 #[test]
2558 fn scheduled_handlers_returns_only_scheduled() {
2559 let mut engine = create_test_engine();
2560 engine.register(EchoWorkflow).unwrap();
2561 engine.register(ScheduledWorkflow::new()).unwrap();
2562 engine.register(FailingWorkflow).unwrap();
2563
2564 let scheduled = engine.scheduled_handlers();
2565 assert_eq!(scheduled.len(), 1);
2566 assert_eq!(scheduled[0].0, "scheduled");
2567 assert_eq!(scheduled[0].1.as_str(), "0 0 * * * *");
2568 }
2569
2570 #[test]
2571 fn scheduled_handlers_empty_when_none_scheduled() {
2572 let mut engine = create_test_engine();
2573 engine.register(EchoWorkflow).unwrap();
2574 engine.register(FailingWorkflow).unwrap();
2575
2576 let scheduled = engine.scheduled_handlers();
2577 assert!(scheduled.is_empty());
2578 }
2579
2580 struct BadCategoryWorkflow(&'static str);
2581
2582 impl WorkflowHandler for BadCategoryWorkflow {
2583 fn name(&self) -> &str {
2584 "bad-category"
2585 }
2586 fn category(&self) -> Option<&str> {
2587 Some(self.0)
2588 }
2589 fn execute<'a>(
2590 &'a self,
2591 _ctx: &'a mut WorkflowContext,
2592 ) -> crate::handler::HandlerFuture<'a> {
2593 Box::pin(async move { Ok(()) })
2594 }
2595 }
2596
2597 #[test]
2598 fn engine_register_rejects_empty_category() {
2599 let mut engine = create_test_engine();
2600 let err = engine.register(BadCategoryWorkflow("")).unwrap_err();
2601 match err {
2602 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty category")),
2603 other => panic!("expected InvalidWorkflow, got {other:?}"),
2604 }
2605 }
2606
2607 #[test]
2608 fn engine_register_rejects_leading_slash_category() {
2609 let mut engine = create_test_engine();
2610 let err = engine
2611 .register(BadCategoryWorkflow("/data/etl"))
2612 .unwrap_err();
2613 match err {
2614 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("leading '/'")),
2615 other => panic!("expected InvalidWorkflow, got {other:?}"),
2616 }
2617 }
2618
2619 #[test]
2620 fn engine_register_rejects_trailing_slash_category() {
2621 let mut engine = create_test_engine();
2622 let err = engine
2623 .register(BadCategoryWorkflow("data/etl/"))
2624 .unwrap_err();
2625 match err {
2626 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("trailing '/'")),
2627 other => panic!("expected InvalidWorkflow, got {other:?}"),
2628 }
2629 }
2630
2631 #[test]
2632 fn engine_register_rejects_double_slash_category() {
2633 let mut engine = create_test_engine();
2634 let err = engine
2635 .register(BadCategoryWorkflow("data//etl"))
2636 .unwrap_err();
2637 match err {
2638 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty segment")),
2639 other => panic!("expected InvalidWorkflow, got {other:?}"),
2640 }
2641 }
2642
2643 #[test]
2644 fn engine_register_rejects_whitespace_only_segment_category() {
2645 let mut engine = create_test_engine();
2646 let err = engine
2647 .register(BadCategoryWorkflow("data/ /etl"))
2648 .unwrap_err();
2649 match err {
2650 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("whitespace-only segment")),
2651 other => panic!("expected InvalidWorkflow, got {other:?}"),
2652 }
2653 }
2654
2655 #[test]
2656 fn engine_register_accepts_valid_nested_category() {
2657 let mut engine = create_test_engine();
2658 assert!(engine.register(CategorizedWorkflow).is_ok());
2659 }
2660
2661 #[tokio::test]
2662 async fn engine_unknown_workflow_returns_error() {
2663 let engine = create_test_engine();
2664 let result = engine
2665 .run_handler("unknown", TriggerKind::Manual, json!({}))
2666 .await;
2667 assert!(result.is_err());
2668 match result {
2669 Err(EngineError::InvalidWorkflow(msg)) => {
2670 assert!(msg.contains("no handler registered"));
2671 }
2672 _ => panic!("expected InvalidWorkflow error"),
2673 }
2674 }
2675
2676 #[tokio::test]
2677 async fn engine_enqueue_handler_creates_pending_run() {
2678 let mut engine = create_test_engine();
2679 engine.register(EchoWorkflow).unwrap();
2680
2681 let run = engine
2682 .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
2683 .await
2684 .unwrap();
2685 assert_eq!(run.status.state, RunStatus::Pending);
2686 assert_eq!(run.workflow_name, "echo-workflow");
2687 }
2688
2689 #[tokio::test]
2690 async fn enqueue_handler_leaves_the_run_unattributed() {
2691 let mut engine = create_test_engine();
2692 engine.register(EchoWorkflow).unwrap();
2693
2694 let run = engine
2695 .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
2696 .await
2697 .unwrap();
2698
2699 assert!(run.created_by.is_none());
2700 }
2701
2702 #[tokio::test]
2703 async fn enqueue_handler_with_options_records_the_author() {
2704 let mut engine = create_test_engine();
2705 engine.register(EchoWorkflow).unwrap();
2706 let actor = RunActor::User {
2707 user_id: Uuid::now_v7(),
2708 };
2709
2710 let run = engine
2711 .enqueue_handler_with_options(
2712 "echo-workflow",
2713 TriggerKind::Api,
2714 json!({}),
2715 EnqueueOptions {
2716 created_by: Some(actor.clone()),
2717 ..Default::default()
2718 },
2719 )
2720 .await
2721 .unwrap()
2722 .into_run();
2723
2724 assert_eq!(run.created_by, Some(actor));
2725 }
2726
2727 #[tokio::test]
2728 async fn enqueue_handler_with_options_accepts_no_author() {
2729 let mut engine = create_test_engine();
2730 engine.register(EchoWorkflow).unwrap();
2731
2732 let run = engine
2733 .enqueue_handler_with_options(
2734 "echo-workflow",
2735 TriggerKind::Cron {
2736 schedule: "0 * * * * *".to_string(),
2737 },
2738 json!({}),
2739 EnqueueOptions::default(),
2740 )
2741 .await
2742 .unwrap()
2743 .into_run();
2744
2745 assert!(run.created_by.is_none());
2746 }
2747
2748 #[tokio::test]
2749 async fn run_handler_leaves_the_run_unattributed() {
2750 let mut engine = create_test_engine();
2751 engine.register(EchoWorkflow).unwrap();
2752
2753 let run = engine
2754 .run_handler("echo-workflow", TriggerKind::Manual, json!({}))
2755 .await
2756 .unwrap()
2757 .run;
2758
2759 assert!(run.created_by.is_none());
2760 }
2761
2762 #[tokio::test]
2763 async fn engine_register_boxed() {
2764 let mut engine = create_test_engine();
2765 let handler: Box<dyn WorkflowHandler> = Box::new(EchoWorkflow);
2766 let result = engine.register_boxed(handler);
2767 assert!(result.is_ok());
2768 assert_eq!(engine.handler_names().len(), 1);
2769 }
2770
2771 #[tokio::test]
2772 async fn engine_store_and_provider_accessors() {
2773 let store = Arc::new(InMemoryStore::new());
2774 let inner = ClaudeCodeProvider::new();
2775 let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
2776 inner,
2777 "/tmp/ironflow-fixtures",
2778 ));
2779 let engine = Engine::new(store.clone(), provider.clone());
2780
2781 let _ = engine.store();
2783 let _ = engine.provider();
2784 }
2785
2786 use crate::operation::{Operation, OperationContext};
2791 use async_trait::async_trait;
2792 use ironflow_core::error::OperationError;
2793 use ironflow_store::models::StepKind;
2794
2795 struct FakeGitlabOp {
2796 project_id: u64,
2797 title: String,
2798 }
2799
2800 #[async_trait]
2801 impl Operation for FakeGitlabOp {
2802 fn kind(&self) -> &str {
2803 "gitlab"
2804 }
2805
2806 async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
2807 Ok(json!({
2808 "issue_id": 42,
2809 "project_id": self.project_id,
2810 "title": self.title,
2811 }))
2812 }
2813
2814 fn input(&self) -> Option<Value> {
2815 Some(json!({
2816 "project_id": self.project_id,
2817 "title": self.title,
2818 }))
2819 }
2820 }
2821
2822 struct FailingOp;
2823
2824 #[async_trait]
2825 impl Operation for FailingOp {
2826 fn kind(&self) -> &str {
2827 "broken-service"
2828 }
2829
2830 async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
2831 Err(OperationError::Http {
2832 status: None,
2833 message: "service unavailable".to_string(),
2834 })
2835 }
2836 }
2837
2838 struct OperationWorkflow;
2839
2840 impl WorkflowHandler for OperationWorkflow {
2841 fn name(&self) -> &str {
2842 "operation-workflow"
2843 }
2844
2845 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2846 Box::pin(async move {
2847 let op = FakeGitlabOp {
2848 project_id: 123,
2849 title: "Bug report".to_string(),
2850 };
2851 ctx.operation("create-issue", &op).await?;
2852 Ok(())
2853 })
2854 }
2855 }
2856
2857 struct FailingOperationWorkflow;
2858
2859 impl WorkflowHandler for FailingOperationWorkflow {
2860 fn name(&self) -> &str {
2861 "failing-operation-workflow"
2862 }
2863
2864 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2865 Box::pin(async move {
2866 ctx.operation("broken-call", &FailingOp).await?;
2867 Ok(())
2868 })
2869 }
2870 }
2871
2872 struct MixedWorkflow;
2873
2874 impl WorkflowHandler for MixedWorkflow {
2875 fn name(&self) -> &str {
2876 "mixed-workflow"
2877 }
2878
2879 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2880 Box::pin(async move {
2881 ctx.shell("build", ShellConfig::new("echo built")).await?;
2882 let op = FakeGitlabOp {
2883 project_id: 456,
2884 title: "Deploy done".to_string(),
2885 };
2886 let result = ctx.operation("notify-gitlab", &op).await?;
2887 assert_eq!(result.output["issue_id"], 42);
2888 Ok(())
2889 })
2890 }
2891 }
2892
2893 #[tokio::test]
2894 async fn operation_step_happy_path() {
2895 let mut engine = create_test_engine();
2896 engine.register(OperationWorkflow).unwrap();
2897
2898 let run = engine
2899 .run_handler("operation-workflow", TriggerKind::Manual, json!({}))
2900 .await
2901 .unwrap()
2902 .run;
2903
2904 assert_eq!(run.status.state, RunStatus::Completed);
2905
2906 let steps = engine.store().list_steps(run.id).await.unwrap();
2907
2908 assert_eq!(steps.len(), 1);
2909 assert_eq!(steps[0].name, "create-issue");
2910 assert_eq!(steps[0].kind, StepKind::Custom("gitlab".to_string()));
2911 assert_eq!(
2912 steps[0].status.state,
2913 ironflow_store::models::StepStatus::Completed
2914 );
2915
2916 let output = steps[0].output.as_ref().unwrap();
2917 assert_eq!(output["issue_id"], 42);
2918 assert_eq!(output["project_id"], 123);
2919
2920 let input = steps[0].input.as_ref().unwrap();
2921 assert_eq!(input["project_id"], 123);
2922 assert_eq!(input["title"], "Bug report");
2923 }
2924
2925 #[tokio::test]
2926 async fn operation_step_failure_marks_run_failed() {
2927 let mut engine = create_test_engine();
2928 engine.register(FailingOperationWorkflow).unwrap();
2929
2930 let result = engine
2931 .run_handler("failing-operation-workflow", TriggerKind::Manual, json!({}))
2932 .await;
2933
2934 assert!(result.is_err());
2935 }
2936
2937 #[tokio::test]
2938 async fn operation_mixed_with_shell_steps() {
2939 let mut engine = create_test_engine();
2940 engine.register(MixedWorkflow).unwrap();
2941
2942 let run = engine
2943 .run_handler("mixed-workflow", TriggerKind::Manual, json!({}))
2944 .await
2945 .unwrap()
2946 .run;
2947
2948 assert_eq!(run.status.state, RunStatus::Completed);
2949
2950 let steps = engine.store().list_steps(run.id).await.unwrap();
2951
2952 assert_eq!(steps.len(), 2);
2953 assert_eq!(steps[0].kind, StepKind::Shell);
2954 assert_eq!(steps[1].kind, StepKind::Custom("gitlab".to_string()));
2955 assert_eq!(steps[0].position, 0);
2956 assert_eq!(steps[1].position, 1);
2957 }
2958
2959 use crate::config::ApprovalConfig;
2964
2965 struct SingleApprovalWorkflow;
2966
2967 impl WorkflowHandler for SingleApprovalWorkflow {
2968 fn name(&self) -> &str {
2969 "single-approval"
2970 }
2971
2972 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2973 Box::pin(async move {
2974 ctx.shell("build", ShellConfig::new("echo built")).await?;
2975 ctx.approval("gate", ApprovalConfig::new("OK?")).await?;
2976 ctx.shell("deploy", ShellConfig::new("echo deployed"))
2977 .await?;
2978 Ok(())
2979 })
2980 }
2981 }
2982
2983 struct DoubleApprovalWorkflow;
2984
2985 impl WorkflowHandler for DoubleApprovalWorkflow {
2986 fn name(&self) -> &str {
2987 "double-approval"
2988 }
2989
2990 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2991 Box::pin(async move {
2992 ctx.shell("build", ShellConfig::new("echo built")).await?;
2993 ctx.approval("staging-gate", ApprovalConfig::new("Deploy staging?"))
2994 .await?;
2995 ctx.shell("deploy-staging", ShellConfig::new("echo staging"))
2996 .await?;
2997 ctx.approval("prod-gate", ApprovalConfig::new("Deploy prod?"))
2998 .await?;
2999 ctx.shell("deploy-prod", ShellConfig::new("echo prod"))
3000 .await?;
3001 Ok(())
3002 })
3003 }
3004 }
3005
3006 #[tokio::test]
3007 async fn approval_pauses_run() {
3008 let mut engine = create_test_engine();
3009 engine.register(SingleApprovalWorkflow).unwrap();
3010
3011 let run = engine
3012 .run_handler("single-approval", TriggerKind::Manual, json!({}))
3013 .await
3014 .unwrap()
3015 .run;
3016
3017 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3018
3019 let steps = engine.store().list_steps(run.id).await.unwrap();
3020 assert_eq!(steps.len(), 2); assert_eq!(steps[0].kind, StepKind::Shell);
3022 assert_eq!(steps[0].status.state, StepStatus::Completed);
3023 assert_eq!(steps[1].kind, StepKind::Approval);
3024 assert_eq!(steps[1].status.state, StepStatus::AwaitingApproval);
3025 }
3026
3027 #[tokio::test]
3028 async fn approval_resume_completes_run() {
3029 let mut engine = create_test_engine();
3030 engine.register(SingleApprovalWorkflow).unwrap();
3031
3032 let run = engine
3034 .run_handler("single-approval", TriggerKind::Manual, json!({}))
3035 .await
3036 .unwrap()
3037 .run;
3038 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3039
3040 engine
3042 .store()
3043 .update_run_status(run.id, RunStatus::Running)
3044 .await
3045 .unwrap();
3046
3047 let resumed = engine.resume_run(run.id).await.unwrap().run;
3049 assert_eq!(resumed.status.state, RunStatus::Completed);
3050
3051 let steps = engine.store().list_steps(run.id).await.unwrap();
3052 assert_eq!(steps.len(), 3); assert_eq!(steps[0].name, "build");
3054 assert_eq!(steps[0].status.state, StepStatus::Completed);
3055 assert_eq!(steps[1].name, "gate");
3056 assert_eq!(steps[1].kind, StepKind::Approval);
3057 assert_eq!(steps[1].status.state, StepStatus::Completed);
3058 assert_eq!(steps[2].name, "deploy");
3059 assert_eq!(steps[2].status.state, StepStatus::Completed);
3060 }
3061
3062 #[tokio::test]
3063 async fn double_approval_two_resumes() {
3064 let mut engine = create_test_engine();
3065 engine.register(DoubleApprovalWorkflow).unwrap();
3066
3067 let run = engine
3069 .run_handler("double-approval", TriggerKind::Manual, json!({}))
3070 .await
3071 .unwrap()
3072 .run;
3073 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3074
3075 let steps = engine.store().list_steps(run.id).await.unwrap();
3076 assert_eq!(steps.len(), 2); engine
3080 .store()
3081 .update_run_status(run.id, RunStatus::Running)
3082 .await
3083 .unwrap();
3084
3085 let resumed = engine.resume_run(run.id).await.unwrap().run;
3086 assert_eq!(resumed.status.state, RunStatus::AwaitingApproval);
3087
3088 let steps = engine.store().list_steps(run.id).await.unwrap();
3089 assert_eq!(steps.len(), 4); engine
3093 .store()
3094 .update_run_status(run.id, RunStatus::Running)
3095 .await
3096 .unwrap();
3097
3098 let final_run = engine.resume_run(run.id).await.unwrap().run;
3099 assert_eq!(final_run.status.state, RunStatus::Completed);
3100
3101 let steps = engine.store().list_steps(run.id).await.unwrap();
3102 assert_eq!(steps.len(), 5);
3103 assert_eq!(steps[0].name, "build");
3104 assert_eq!(steps[1].name, "staging-gate");
3105 assert_eq!(steps[2].name, "deploy-staging");
3106 assert_eq!(steps[3].name, "prod-gate");
3107 assert_eq!(steps[4].name, "deploy-prod");
3108
3109 for step in &steps {
3110 assert_eq!(step.status.state, StepStatus::Completed);
3111 }
3112 }
3113
3114 use ironflow_store::models::{NewStep, StepUpdate, step_trace_id};
3119
3120 async fn create_step_with_status(
3121 store: &Arc<dyn Store>,
3122 run_id: Uuid,
3123 name: &str,
3124 position: u32,
3125 status: StepStatus,
3126 ) -> ironflow_store::models::Step {
3127 let step = store
3128 .create_step(NewStep {
3129 run_id,
3130 trace_id: step_trace_id(run_id, name, position),
3131 name: name.to_string(),
3132 kind: StepKind::Shell,
3133 position,
3134 input: None,
3135 is_error_handler: false,
3136 })
3137 .await
3138 .unwrap();
3139
3140 match status {
3141 StepStatus::Pending => {}
3142 StepStatus::Running => {
3143 store
3144 .update_step(
3145 step.id,
3146 StepUpdate {
3147 status: Some(StepStatus::Running),
3148 ..StepUpdate::default()
3149 },
3150 )
3151 .await
3152 .unwrap();
3153 }
3154 StepStatus::Completed => {
3155 store
3156 .update_step(
3157 step.id,
3158 StepUpdate {
3159 status: Some(StepStatus::Running),
3160 ..StepUpdate::default()
3161 },
3162 )
3163 .await
3164 .unwrap();
3165 store
3166 .update_step(
3167 step.id,
3168 StepUpdate {
3169 status: Some(StepStatus::Completed),
3170 ..StepUpdate::default()
3171 },
3172 )
3173 .await
3174 .unwrap();
3175 }
3176 StepStatus::AwaitingApproval => {
3177 store
3178 .update_step(
3179 step.id,
3180 StepUpdate {
3181 status: Some(StepStatus::Running),
3182 ..StepUpdate::default()
3183 },
3184 )
3185 .await
3186 .unwrap();
3187 store
3188 .update_step(
3189 step.id,
3190 StepUpdate {
3191 status: Some(StepStatus::AwaitingApproval),
3192 ..StepUpdate::default()
3193 },
3194 )
3195 .await
3196 .unwrap();
3197 }
3198 _ => panic!("unsupported status for test helper: {status}"),
3199 }
3200
3201 store.get_step(step.id).await.unwrap().unwrap()
3202 }
3203
3204 #[tokio::test]
3205 async fn fail_orphaned_steps_marks_running_as_failed() {
3206 let engine = create_test_engine();
3207 let run = engine
3208 .store()
3209 .create_run(NewRun {
3210 created_by: None,
3211 workflow_name: "test".to_string(),
3212 trigger: TriggerKind::Manual,
3213 payload: json!({}),
3214 max_retries: 0,
3215 handler_version: None,
3216 labels: HashMap::new(),
3217 scheduled_at: None,
3218 idempotency_key: None,
3219 max_cost_usd: None,
3220 })
3221 .await
3222 .unwrap()
3223 .into_run();
3224
3225 let step = create_step_with_status(
3226 engine.store(),
3227 run.id,
3228 "running-step",
3229 0,
3230 StepStatus::Running,
3231 )
3232 .await;
3233
3234 engine
3235 .fail_orphaned_steps(run.id, "parent run timed out")
3236 .await
3237 .unwrap();
3238
3239 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
3240 assert_eq!(updated.status.state, StepStatus::Failed);
3241 assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
3242 assert!(updated.completed_at.is_some());
3243 }
3244
3245 #[tokio::test]
3246 async fn fail_orphaned_steps_marks_pending_as_skipped() {
3247 let engine = create_test_engine();
3248 let run = engine
3249 .store()
3250 .create_run(NewRun {
3251 created_by: None,
3252 workflow_name: "test".to_string(),
3253 trigger: TriggerKind::Manual,
3254 payload: json!({}),
3255 max_retries: 0,
3256 handler_version: None,
3257 labels: HashMap::new(),
3258 scheduled_at: None,
3259 idempotency_key: None,
3260 max_cost_usd: None,
3261 })
3262 .await
3263 .unwrap()
3264 .into_run();
3265
3266 let step = create_step_with_status(
3267 engine.store(),
3268 run.id,
3269 "pending-step",
3270 0,
3271 StepStatus::Pending,
3272 )
3273 .await;
3274
3275 engine
3276 .fail_orphaned_steps(run.id, "parent run timed out")
3277 .await
3278 .unwrap();
3279
3280 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
3281 assert_eq!(updated.status.state, StepStatus::Skipped);
3282 assert!(updated.error.is_none());
3283 assert!(updated.completed_at.is_some());
3284 }
3285
3286 #[tokio::test]
3287 async fn fail_orphaned_steps_marks_awaiting_approval_as_failed() {
3288 let engine = create_test_engine();
3289 let run = engine
3290 .store()
3291 .create_run(NewRun {
3292 created_by: None,
3293 workflow_name: "test".to_string(),
3294 trigger: TriggerKind::Manual,
3295 payload: json!({}),
3296 max_retries: 0,
3297 handler_version: None,
3298 labels: HashMap::new(),
3299 scheduled_at: None,
3300 idempotency_key: None,
3301 max_cost_usd: None,
3302 })
3303 .await
3304 .unwrap()
3305 .into_run();
3306
3307 let step = create_step_with_status(
3308 engine.store(),
3309 run.id,
3310 "approval-step",
3311 0,
3312 StepStatus::AwaitingApproval,
3313 )
3314 .await;
3315
3316 engine
3317 .fail_orphaned_steps(run.id, "parent run timed out")
3318 .await
3319 .unwrap();
3320
3321 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
3322 assert_eq!(updated.status.state, StepStatus::Failed);
3323 assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
3324 assert!(updated.completed_at.is_some());
3325 }
3326
3327 #[tokio::test]
3328 async fn fail_orphaned_steps_skips_terminal_steps() {
3329 let engine = create_test_engine();
3330 let run = engine
3331 .store()
3332 .create_run(NewRun {
3333 created_by: None,
3334 workflow_name: "test".to_string(),
3335 trigger: TriggerKind::Manual,
3336 payload: json!({}),
3337 max_retries: 0,
3338 handler_version: None,
3339 labels: HashMap::new(),
3340 scheduled_at: None,
3341 idempotency_key: None,
3342 max_cost_usd: None,
3343 })
3344 .await
3345 .unwrap()
3346 .into_run();
3347
3348 let completed_step =
3349 create_step_with_status(engine.store(), run.id, "done", 0, StepStatus::Completed).await;
3350 let running_step =
3351 create_step_with_status(engine.store(), run.id, "in-flight", 1, StepStatus::Running)
3352 .await;
3353
3354 engine
3355 .fail_orphaned_steps(run.id, "parent run timed out")
3356 .await
3357 .unwrap();
3358
3359 let completed = engine
3360 .store()
3361 .get_step(completed_step.id)
3362 .await
3363 .unwrap()
3364 .unwrap();
3365 assert_eq!(completed.status.state, StepStatus::Completed);
3366
3367 let failed = engine
3368 .store()
3369 .get_step(running_step.id)
3370 .await
3371 .unwrap()
3372 .unwrap();
3373 assert_eq!(failed.status.state, StepStatus::Failed);
3374 }
3375
3376 #[tokio::test]
3377 async fn fail_orphaned_steps_mixed_states() {
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 max_cost_usd: None,
3392 })
3393 .await
3394 .unwrap()
3395 .into_run();
3396
3397 let s_completed =
3398 create_step_with_status(engine.store(), run.id, "step-1", 0, StepStatus::Completed)
3399 .await;
3400 let s_running =
3401 create_step_with_status(engine.store(), run.id, "step-2", 1, StepStatus::Running).await;
3402 let s_pending =
3403 create_step_with_status(engine.store(), run.id, "step-3", 2, StepStatus::Pending).await;
3404
3405 engine.fail_orphaned_steps(run.id, "timeout").await.unwrap();
3406
3407 let r_completed = engine
3408 .store()
3409 .get_step(s_completed.id)
3410 .await
3411 .unwrap()
3412 .unwrap();
3413 assert_eq!(r_completed.status.state, StepStatus::Completed);
3414
3415 let r_running = engine
3416 .store()
3417 .get_step(s_running.id)
3418 .await
3419 .unwrap()
3420 .unwrap();
3421 assert_eq!(r_running.status.state, StepStatus::Failed);
3422 assert_eq!(r_running.error.as_deref(), Some("timeout"));
3423
3424 let r_pending = engine
3425 .store()
3426 .get_step(s_pending.id)
3427 .await
3428 .unwrap()
3429 .unwrap();
3430 assert_eq!(r_pending.status.state, StepStatus::Skipped);
3431 assert!(r_pending.error.is_none());
3432 }
3433
3434 #[tokio::test]
3435 async fn fail_orphaned_steps_no_steps_is_noop() {
3436 let engine = create_test_engine();
3437 let run = engine
3438 .store()
3439 .create_run(NewRun {
3440 created_by: None,
3441 workflow_name: "test".to_string(),
3442 trigger: TriggerKind::Manual,
3443 payload: json!({}),
3444 max_retries: 0,
3445 handler_version: None,
3446 labels: HashMap::new(),
3447 scheduled_at: None,
3448 idempotency_key: None,
3449 max_cost_usd: None,
3450 })
3451 .await
3452 .unwrap()
3453 .into_run();
3454
3455 let result = engine.fail_orphaned_steps(run.id, "timeout").await;
3456 assert!(result.is_ok());
3457 }
3458
3459 #[tokio::test]
3460 async fn fail_orphaned_steps_preserves_existing_error() {
3461 let engine = create_test_engine();
3462 let run = engine
3463 .store()
3464 .create_run(NewRun {
3465 created_by: None,
3466 workflow_name: "test".to_string(),
3467 trigger: TriggerKind::Manual,
3468 payload: json!({}),
3469 max_retries: 0,
3470 handler_version: None,
3471 labels: HashMap::new(),
3472 scheduled_at: None,
3473 idempotency_key: None,
3474 max_cost_usd: None,
3475 })
3476 .await
3477 .unwrap()
3478 .into_run();
3479
3480 let step_with_error = create_step_with_status(
3481 engine.store(),
3482 run.id,
3483 "already-errored",
3484 0,
3485 StepStatus::Running,
3486 )
3487 .await;
3488
3489 engine
3490 .store()
3491 .update_step(
3492 step_with_error.id,
3493 StepUpdate {
3494 error: Some("real error from provider".to_string()),
3495 ..StepUpdate::default()
3496 },
3497 )
3498 .await
3499 .unwrap();
3500
3501 let step_no_error = create_step_with_status(
3502 engine.store(),
3503 run.id,
3504 "no-error-yet",
3505 1,
3506 StepStatus::Running,
3507 )
3508 .await;
3509
3510 engine
3511 .fail_orphaned_steps(run.id, "parent run failed")
3512 .await
3513 .unwrap();
3514
3515 let updated_with = engine
3516 .store()
3517 .get_step(step_with_error.id)
3518 .await
3519 .unwrap()
3520 .unwrap();
3521 assert_eq!(updated_with.status.state, StepStatus::Failed);
3522 assert_eq!(
3523 updated_with.error.as_deref(),
3524 Some("real error from provider"),
3525 );
3526
3527 let updated_without = engine
3528 .store()
3529 .get_step(step_no_error.id)
3530 .await
3531 .unwrap()
3532 .unwrap();
3533 assert_eq!(updated_without.status.state, StepStatus::Failed);
3534 assert_eq!(updated_without.error.as_deref(), Some("parent run failed"),);
3535 }
3536}