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