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 ..RunUpdate::default()
1847 },
1848 )
1849 .await?;
1850
1851 info!(
1852 run_id = %run_id,
1853 status = %final_status,
1854 cost_usd = %ctx.total_cost_usd(),
1855 duration_ms = total_duration,
1856 "run completed"
1857 );
1858 }
1859 Err(EngineError::ApprovalRequired {
1860 run_id: approval_run_id,
1861 step_id,
1862 ref message,
1863 }) => {
1864 final_status = RunStatus::AwaitingApproval;
1865 final_run = self
1866 .store
1867 .update_run_returning(
1868 run_id,
1869 RunUpdate {
1870 status: Some(RunStatus::AwaitingApproval),
1871 cost_usd: Some(ctx.total_cost_usd()),
1872 duration_ms: Some(total_duration),
1873 ..RunUpdate::default()
1874 },
1875 )
1876 .await?;
1877
1878 info!(
1879 run_id = %approval_run_id,
1880 step_id = %step_id,
1881 message = %message,
1882 "run awaiting approval"
1883 );
1884
1885 self.publish_approval_requested(approval_run_id, step_id, message)
1886 .await?;
1887 }
1888 Err(EngineError::ChildSuspended {
1889 run_id: child_run_id,
1890 ref cause,
1891 }) => {
1892 final_status = cause.suspension_status();
1893 final_run = self
1897 .store
1898 .update_run_returning(
1899 run_id,
1900 RunUpdate {
1901 status: Some(final_status),
1902 cost_usd: Some(ctx.total_cost_usd()),
1903 duration_ms: Some(total_duration),
1904 ..RunUpdate::default()
1905 },
1906 )
1907 .await?;
1908
1909 let leaf = cause.suspension_leaf();
1910 info!(
1911 run_id = %run_id,
1912 child_run_id = %child_run_id,
1913 status = %final_status,
1914 cause = %leaf,
1915 "run suspended with its child run"
1916 );
1917
1918 match leaf {
1919 EngineError::ApprovalRequired {
1920 run_id: approval_run_id,
1921 step_id,
1922 message,
1923 } => {
1924 self.publish_approval_requested(*approval_run_id, *step_id, message)
1925 .await?;
1926 }
1927 EngineError::SignalWaiting {
1928 run_id: wait_run_id,
1929 step_id,
1930 step_name,
1931 name,
1932 key,
1933 deadline_at,
1934 } => {
1935 self.event_publisher
1936 .publish(Event::SignalAwaited(SignalAwaitedEvent {
1937 run_id: *wait_run_id,
1938 step_id: *step_id,
1939 step_name: step_name.clone(),
1940 name: name.clone(),
1941 key: key.clone(),
1942 deadline_at: *deadline_at,
1943 at: Utc::now(),
1944 }));
1945 }
1946 _ => {}
1949 }
1950 }
1951 Err(EngineError::HumanInputRequired {
1952 run_id: input_run_id,
1953 step_id,
1954 ref message,
1955 }) => {
1956 final_status = RunStatus::AwaitingApproval;
1957 final_run = self
1958 .store
1959 .update_run_returning(
1960 run_id,
1961 RunUpdate {
1962 status: Some(RunStatus::AwaitingApproval),
1963 cost_usd: Some(ctx.total_cost_usd()),
1964 duration_ms: Some(total_duration),
1965 ..RunUpdate::default()
1966 },
1967 )
1968 .await?;
1969
1970 info!(
1972 run_id = %input_run_id,
1973 step_id = %step_id,
1974 message = %message,
1975 "run awaiting human input"
1976 );
1977 }
1978 Err(EngineError::DelaySleeping {
1979 run_id: delay_run_id,
1980 step_id,
1981 wake_at,
1982 }) => {
1983 final_status = RunStatus::Sleeping;
1984 final_run = self
1985 .store
1986 .update_run_returning(
1987 run_id,
1988 RunUpdate {
1989 status: Some(RunStatus::Sleeping),
1990 cost_usd: Some(ctx.total_cost_usd()),
1991 duration_ms: Some(total_duration),
1992 scheduled_at: Some(wake_at),
1993 ..RunUpdate::default()
1994 },
1995 )
1996 .await?;
1997
1998 info!(
1999 run_id = %delay_run_id,
2000 step_id = %step_id,
2001 wake_at = %wake_at,
2002 "run sleeping until delay elapses"
2003 );
2004 }
2005 Err(EngineError::SignalWaiting {
2006 run_id: wait_run_id,
2007 step_id,
2008 ref step_name,
2009 ref name,
2010 ref key,
2011 deadline_at,
2012 }) => {
2013 final_status = RunStatus::Sleeping;
2014 let waiting = self
2018 .store
2019 .suspend_run_on_signal(run_id, step_id, deadline_at)
2020 .await?;
2021 final_run = self
2022 .store
2023 .update_run_returning(
2024 run_id,
2025 RunUpdate {
2026 cost_usd: Some(ctx.total_cost_usd()),
2027 duration_ms: Some(total_duration),
2028 ..RunUpdate::default()
2029 },
2030 )
2031 .await?;
2032
2033 if waiting {
2034 self.event_publisher
2035 .publish(Event::SignalAwaited(SignalAwaitedEvent {
2036 run_id: wait_run_id,
2037 step_id,
2038 step_name: step_name.clone(),
2039 name: name.clone(),
2040 key: key.clone(),
2041 deadline_at,
2042 at: Utc::now(),
2043 }));
2044 }
2045
2046 info!(
2047 run_id = %wait_run_id,
2048 step_id = %step_id,
2049 signal = %name,
2050 key = %key,
2051 deadline_at = %deadline_at,
2052 waiting,
2053 "run sleeping until a signal arrives"
2054 );
2055 }
2056 Err(err) => {
2057 let guardrail_stop = matches!(
2061 err,
2062 EngineError::RunBudgetExceeded { .. } | EngineError::WorkflowGuardRejected(_)
2063 );
2064
2065 final_status = if guardrail_stop {
2066 if let Err(store_err) = self
2067 .store
2068 .update_run(
2069 run_id,
2070 RunUpdate {
2071 status: Some(RunStatus::Cancelled),
2072 error: Some(err.to_string()),
2073 cost_usd: Some(ctx.total_cost_usd()),
2074 duration_ms: Some(total_duration),
2075 completed_at: Some(completed_at),
2076 ..RunUpdate::default()
2077 },
2078 )
2079 .await
2080 {
2081 error!(run_id = %run_id, store_error = %store_err, "failed to persist run cancellation");
2082 }
2083 if let Err(cleanup_err) = self
2084 .fail_orphaned_steps(run_id, "run stopped: guardrail limit reached")
2085 .await
2086 {
2087 error!(run_id = %run_id, store_error = %cleanup_err, "failed to cleanup orphaned steps");
2088 }
2089 RunStatus::Cancelled
2090 } else {
2091 self.fail_or_schedule_retry(
2092 run_id,
2093 &err.to_string(),
2094 is_run_retryable(&err),
2095 Some(ctx.total_cost_usd()),
2096 Some(total_duration),
2097 )
2098 .await
2099 .unwrap_or_else(|store_err| {
2100 error!(run_id = %run_id, store_error = %store_err, "failed to persist run failure");
2101 RunStatus::Failed
2102 })
2103 };
2104
2105 if matches!(err, EngineError::RunBudgetExceeded { .. }) {
2106 self.on_run_budget_exceeded(workflow_name, run_id, &err);
2107 }
2108
2109 error!(run_id = %run_id, status = %final_status, error = %err, "run stopped");
2110
2111 self.publish_run_status_changed(
2112 workflow_name,
2113 run_id,
2114 final_status,
2115 Some(err.to_string()),
2116 ctx,
2117 total_duration,
2118 run_labels,
2119 );
2120
2121 #[cfg(feature = "prometheus")]
2122 self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
2123
2124 return Err(err);
2125 }
2126 }
2127
2128 self.publish_run_status_changed(
2129 workflow_name,
2130 run_id,
2131 final_status,
2132 None,
2133 ctx,
2134 total_duration,
2135 run_labels,
2136 );
2137
2138 #[cfg(feature = "prometheus")]
2139 self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
2140
2141 Ok(WorkflowResult {
2142 run: final_run,
2143 steps: ctx.step_results().to_vec(),
2144 })
2145 }
2146
2147 async fn publish_approval_requested(
2150 &self,
2151 run_id: Uuid,
2152 step_id: Uuid,
2153 message: &str,
2154 ) -> Result<(), EngineError> {
2155 let requirement = self
2156 .store
2157 .get_step(step_id)
2158 .await?
2159 .and_then(|s| s.approval_requirement);
2160 self.event_publisher
2161 .publish(Event::ApprovalRequested(ApprovalRequestedEvent {
2162 run_id,
2163 step_id,
2164 message: message.to_string(),
2165 requirement,
2166 at: Utc::now(),
2167 }));
2168 Ok(())
2169 }
2170
2171 pub async fn fail_ancestors(&self, run_id: Uuid, reason: &str) -> Result<(), EngineError> {
2200 let mut current = self
2201 .store
2202 .get_run(run_id)
2203 .await?
2204 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
2205 let mut visited = HashSet::from([run_id]);
2207
2208 while let Some(parent_id) = chain_parent(¤t) {
2209 if !visited.insert(parent_id) {
2210 break;
2211 }
2212 let status = self
2213 .fail_or_schedule_retry(parent_id, reason, false, None, None)
2214 .await?;
2215 info!(
2216 run_id = %run_id,
2217 ancestor_run_id = %parent_id,
2218 status = %status,
2219 "ancestor run failed with its child"
2220 );
2221 current = self
2222 .store
2223 .get_run(parent_id)
2224 .await?
2225 .ok_or(EngineError::Store(StoreError::RunNotFound(parent_id)))?;
2226 }
2227
2228 Ok(())
2229 }
2230
2231 #[cfg(feature = "prometheus")]
2233 fn emit_run_metrics(
2234 &self,
2235 workflow_name: &str,
2236 status: RunStatus,
2237 duration_ms: u64,
2238 ctx: &WorkflowContext,
2239 ) {
2240 let status_str = status.to_string();
2241 let wf = workflow_name.to_string();
2242
2243 counter!(RUNS_TOTAL, "workflow" => wf.clone(), "status" => status_str.clone()).increment(1);
2244 histogram!(RUN_DURATION_SECONDS, "workflow" => wf.clone(), "status" => status_str)
2245 .record(duration_ms as f64 / 1000.0);
2246 histogram!(RUN_COST_USD, "workflow" => wf.clone()).record(
2247 ctx.total_cost_usd()
2248 .to_string()
2249 .parse::<f64>()
2250 .unwrap_or(0.0),
2251 );
2252 gauge!(RUNS_ACTIVE, "workflow" => wf).decrement(1.0);
2253 }
2254
2255 fn on_run_budget_exceeded(&self, workflow_name: &str, run_id: Uuid, err: &EngineError) {
2261 let EngineError::RunBudgetExceeded {
2262 limit_usd,
2263 spent_usd,
2264 step_budget_usd,
2265 ..
2266 } = err
2267 else {
2268 return;
2269 };
2270
2271 #[cfg(feature = "prometheus")]
2272 counter!(
2273 RUN_BUDGET_EXCEEDED_TOTAL,
2274 "workflow" => workflow_name.to_string(),
2275 "scope" => "run",
2276 )
2277 .increment(1);
2278
2279 self.event_publisher
2280 .publish(Event::RunBudgetExceeded(RunBudgetExceededEvent {
2281 run_id,
2282 workflow_name: workflow_name.to_string(),
2283 limit_usd: *limit_usd,
2284 spent_usd: *spent_usd,
2285 step_budget_usd: *step_budget_usd,
2286 at: Utc::now(),
2287 }));
2288 }
2289
2290 #[allow(clippy::too_many_arguments)]
2295 fn publish_run_status_changed(
2296 &self,
2297 workflow_name: &str,
2298 run_id: Uuid,
2299 to: RunStatus,
2300 error: Option<String>,
2301 ctx: &WorkflowContext,
2302 duration_ms: u64,
2303 labels: HashMap<String, String>,
2304 ) {
2305 let now = Utc::now();
2306 let cost_usd = ctx.total_cost_usd();
2307 let wf = workflow_name.to_string();
2308
2309 self.event_publisher
2310 .publish(Event::RunStatusChanged(RunStatusChangedEvent {
2311 run_id,
2312 workflow_name: wf.clone(),
2313 from: RunStatus::Running,
2314 to,
2315 error: error.clone(),
2316 cost_usd,
2317 duration_ms,
2318 labels: labels.clone(),
2319 at: now,
2320 }));
2321
2322 if to == RunStatus::Failed {
2323 self.event_publisher
2324 .publish(Event::RunFailed(RunFailedEvent {
2325 run_id,
2326 workflow_name: wf,
2327 error,
2328 cost_usd,
2329 duration_ms,
2330 labels,
2331 at: now,
2332 }));
2333 }
2334 }
2335}
2336
2337impl fmt::Debug for Engine {
2338 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
2339 f.debug_struct("Engine")
2340 .field("handlers", &self.handlers.keys().collect::<Vec<_>>())
2341 .finish_non_exhaustive()
2342 }
2343}
2344
2345#[cfg(test)]
2346mod tests {
2347 use super::*;
2348 use crate::config::ShellConfig;
2349 use crate::handler::{HandlerFuture, WorkflowHandler};
2350 use ironflow_core::providers::claude::ClaudeCodeProvider;
2351 use ironflow_core::providers::record_replay::RecordReplayProvider;
2352 use ironflow_store::memory::InMemoryStore;
2353 use ironflow_store::models::StepStatus;
2354 use serde_json::json;
2355
2356 struct EchoWorkflow;
2358
2359 impl WorkflowHandler for EchoWorkflow {
2360 fn name(&self) -> &str {
2361 "echo-workflow"
2362 }
2363
2364 fn describe(&self) -> WorkflowInfo {
2365 WorkflowInfo {
2366 description: "A simple workflow that echoes hello".to_string(),
2367 source_code: None,
2368 sub_workflows: Vec::new(),
2369 category: None,
2370 version: self.version().map(str::to_string),
2371 compatible_versions: Vec::new(),
2372 input_schema: None,
2373 default_labels: HashMap::new(),
2374 schedule: self.schedule().cloned(),
2375 default_max_cost_usd: self.default_max_cost_usd(),
2376 }
2377 }
2378
2379 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2380 Box::pin(async move {
2381 ctx.shell("greet", ShellConfig::new("echo hello")).await?;
2382 Ok(())
2383 })
2384 }
2385 }
2386
2387 struct FailingWorkflow;
2389
2390 impl WorkflowHandler for FailingWorkflow {
2391 fn name(&self) -> &str {
2392 "failing-workflow"
2393 }
2394
2395 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2396 Box::pin(async move {
2397 ctx.shell("fail", ShellConfig::new("exit 1")).await?;
2398 Ok(())
2399 })
2400 }
2401 }
2402
2403 fn create_test_engine() -> Engine {
2404 let store = Arc::new(InMemoryStore::new());
2405 let inner = ClaudeCodeProvider::new();
2406 let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
2407 inner,
2408 "/tmp/ironflow-fixtures",
2409 ));
2410 Engine::new(store, provider)
2411 }
2412
2413 #[test]
2414 fn engine_new_creates_instance() {
2415 let engine = create_test_engine();
2416 assert_eq!(engine.handler_names().len(), 0);
2417 }
2418
2419 #[test]
2420 fn execution_mode_defaults_to_local() {
2421 let engine = create_test_engine();
2422 assert_eq!(engine.execution_mode(), ExecutionMode::Local);
2423 }
2424
2425 #[test]
2426 fn with_execution_mode_overrides_the_default() {
2427 let engine = create_test_engine().with_execution_mode(ExecutionMode::Workers);
2428 assert_eq!(engine.execution_mode(), ExecutionMode::Workers);
2429 }
2430
2431 #[test]
2432 fn engine_register_handler() {
2433 let mut engine = create_test_engine();
2434 let result = engine.register(EchoWorkflow);
2435 assert!(result.is_ok());
2436 assert_eq!(engine.handler_names().len(), 1);
2437 assert!(engine.handler_names().contains(&"echo-workflow"));
2438 }
2439
2440 #[test]
2441 fn engine_register_duplicate_returns_error() {
2442 let mut engine = create_test_engine();
2443 engine.register(EchoWorkflow).unwrap();
2444 let result = engine.register(EchoWorkflow);
2445 assert!(result.is_err());
2446 }
2447
2448 #[test]
2449 fn engine_get_handler_found() {
2450 let mut engine = create_test_engine();
2451 engine.register(EchoWorkflow).unwrap();
2452 let handler = engine.get_handler("echo-workflow");
2453 assert!(handler.is_some());
2454 }
2455
2456 #[test]
2457 fn engine_get_handler_not_found() {
2458 let engine = create_test_engine();
2459 let handler = engine.get_handler("nonexistent");
2460 assert!(handler.is_none());
2461 }
2462
2463 #[test]
2464 fn engine_handler_names_lists_all() {
2465 let mut engine = create_test_engine();
2466 engine.register(EchoWorkflow).unwrap();
2467 engine.register(FailingWorkflow).unwrap();
2468 let names = engine.handler_names();
2469 assert_eq!(names.len(), 2);
2470 assert!(names.contains(&"echo-workflow"));
2471 assert!(names.contains(&"failing-workflow"));
2472 }
2473
2474 #[test]
2475 fn engine_handler_info_returns_description() {
2476 let mut engine = create_test_engine();
2477 engine.register(EchoWorkflow).unwrap();
2478 let info = engine.handler_info("echo-workflow");
2479 assert!(info.is_some());
2480 let info = info.unwrap();
2481 assert_eq!(info.description, "A simple workflow that echoes hello");
2482 }
2483
2484 struct CategorizedWorkflow;
2485
2486 impl WorkflowHandler for CategorizedWorkflow {
2487 fn name(&self) -> &str {
2488 "categorized"
2489 }
2490 fn category(&self) -> Option<&str> {
2491 Some("data/etl")
2492 }
2493 fn execute<'a>(
2494 &'a self,
2495 _ctx: &'a mut WorkflowContext,
2496 ) -> crate::handler::HandlerFuture<'a> {
2497 Box::pin(async move { Ok(()) })
2498 }
2499 }
2500
2501 #[test]
2502 fn engine_default_describe_propagates_category() {
2503 let mut engine = create_test_engine();
2504 engine.register(CategorizedWorkflow).unwrap();
2505 let info = engine.handler_info("categorized").unwrap();
2506 assert_eq!(info.category.as_deref(), Some("data/etl"));
2507 }
2508
2509 #[test]
2510 fn engine_default_describe_without_category() {
2511 let mut engine = create_test_engine();
2512 engine.register(EchoWorkflow).unwrap();
2513 let info = engine.handler_info("echo-workflow").unwrap();
2514 assert!(info.category.is_none());
2515 }
2516
2517 struct ScheduledWorkflow {
2522 schedule: CronSchedule,
2523 }
2524
2525 impl ScheduledWorkflow {
2526 fn new() -> Self {
2527 Self {
2528 schedule: CronSchedule::new("0 0 * * * *").unwrap(),
2529 }
2530 }
2531 }
2532
2533 impl WorkflowHandler for ScheduledWorkflow {
2534 fn name(&self) -> &str {
2535 "scheduled"
2536 }
2537 fn schedule(&self) -> Option<&CronSchedule> {
2538 Some(&self.schedule)
2539 }
2540 fn execute<'a>(
2541 &'a self,
2542 _ctx: &'a mut WorkflowContext,
2543 ) -> crate::handler::HandlerFuture<'a> {
2544 Box::pin(async move { Ok(()) })
2545 }
2546 }
2547
2548 #[test]
2549 fn engine_default_describe_propagates_schedule() {
2550 let mut engine = create_test_engine();
2551 engine.register(ScheduledWorkflow::new()).unwrap();
2552 let info = engine.handler_info("scheduled").unwrap();
2553 assert_eq!(
2554 info.schedule.as_ref().map(|s| s.as_str()),
2555 Some("0 0 * * * *")
2556 );
2557 }
2558
2559 #[test]
2560 fn engine_default_describe_without_schedule() {
2561 let mut engine = create_test_engine();
2562 engine.register(EchoWorkflow).unwrap();
2563 let info = engine.handler_info("echo-workflow").unwrap();
2564 assert!(info.schedule.is_none());
2565 }
2566
2567 #[test]
2568 fn scheduled_handlers_returns_only_scheduled() {
2569 let mut engine = create_test_engine();
2570 engine.register(EchoWorkflow).unwrap();
2571 engine.register(ScheduledWorkflow::new()).unwrap();
2572 engine.register(FailingWorkflow).unwrap();
2573
2574 let scheduled = engine.scheduled_handlers();
2575 assert_eq!(scheduled.len(), 1);
2576 assert_eq!(scheduled[0].0, "scheduled");
2577 assert_eq!(scheduled[0].1.as_str(), "0 0 * * * *");
2578 }
2579
2580 #[test]
2581 fn scheduled_handlers_empty_when_none_scheduled() {
2582 let mut engine = create_test_engine();
2583 engine.register(EchoWorkflow).unwrap();
2584 engine.register(FailingWorkflow).unwrap();
2585
2586 let scheduled = engine.scheduled_handlers();
2587 assert!(scheduled.is_empty());
2588 }
2589
2590 struct BadCategoryWorkflow(&'static str);
2591
2592 impl WorkflowHandler for BadCategoryWorkflow {
2593 fn name(&self) -> &str {
2594 "bad-category"
2595 }
2596 fn category(&self) -> Option<&str> {
2597 Some(self.0)
2598 }
2599 fn execute<'a>(
2600 &'a self,
2601 _ctx: &'a mut WorkflowContext,
2602 ) -> crate::handler::HandlerFuture<'a> {
2603 Box::pin(async move { Ok(()) })
2604 }
2605 }
2606
2607 #[test]
2608 fn engine_register_rejects_empty_category() {
2609 let mut engine = create_test_engine();
2610 let err = engine.register(BadCategoryWorkflow("")).unwrap_err();
2611 match err {
2612 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty category")),
2613 other => panic!("expected InvalidWorkflow, got {other:?}"),
2614 }
2615 }
2616
2617 #[test]
2618 fn engine_register_rejects_leading_slash_category() {
2619 let mut engine = create_test_engine();
2620 let err = engine
2621 .register(BadCategoryWorkflow("/data/etl"))
2622 .unwrap_err();
2623 match err {
2624 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("leading '/'")),
2625 other => panic!("expected InvalidWorkflow, got {other:?}"),
2626 }
2627 }
2628
2629 #[test]
2630 fn engine_register_rejects_trailing_slash_category() {
2631 let mut engine = create_test_engine();
2632 let err = engine
2633 .register(BadCategoryWorkflow("data/etl/"))
2634 .unwrap_err();
2635 match err {
2636 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("trailing '/'")),
2637 other => panic!("expected InvalidWorkflow, got {other:?}"),
2638 }
2639 }
2640
2641 #[test]
2642 fn engine_register_rejects_double_slash_category() {
2643 let mut engine = create_test_engine();
2644 let err = engine
2645 .register(BadCategoryWorkflow("data//etl"))
2646 .unwrap_err();
2647 match err {
2648 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty segment")),
2649 other => panic!("expected InvalidWorkflow, got {other:?}"),
2650 }
2651 }
2652
2653 #[test]
2654 fn engine_register_rejects_whitespace_only_segment_category() {
2655 let mut engine = create_test_engine();
2656 let err = engine
2657 .register(BadCategoryWorkflow("data/ /etl"))
2658 .unwrap_err();
2659 match err {
2660 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("whitespace-only segment")),
2661 other => panic!("expected InvalidWorkflow, got {other:?}"),
2662 }
2663 }
2664
2665 #[test]
2666 fn engine_register_accepts_valid_nested_category() {
2667 let mut engine = create_test_engine();
2668 assert!(engine.register(CategorizedWorkflow).is_ok());
2669 }
2670
2671 #[tokio::test]
2672 async fn engine_unknown_workflow_returns_error() {
2673 let engine = create_test_engine();
2674 let result = engine
2675 .run_handler("unknown", TriggerKind::Manual, json!({}))
2676 .await;
2677 assert!(result.is_err());
2678 match result {
2679 Err(EngineError::InvalidWorkflow(msg)) => {
2680 assert!(msg.contains("no handler registered"));
2681 }
2682 _ => panic!("expected InvalidWorkflow error"),
2683 }
2684 }
2685
2686 #[tokio::test]
2687 async fn engine_enqueue_handler_creates_pending_run() {
2688 let mut engine = create_test_engine();
2689 engine.register(EchoWorkflow).unwrap();
2690
2691 let run = engine
2692 .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
2693 .await
2694 .unwrap();
2695 assert_eq!(run.status.state, RunStatus::Pending);
2696 assert_eq!(run.workflow_name, "echo-workflow");
2697 }
2698
2699 #[tokio::test]
2700 async fn enqueue_handler_leaves_the_run_unattributed() {
2701 let mut engine = create_test_engine();
2702 engine.register(EchoWorkflow).unwrap();
2703
2704 let run = engine
2705 .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
2706 .await
2707 .unwrap();
2708
2709 assert!(run.created_by.is_none());
2710 }
2711
2712 #[tokio::test]
2713 async fn enqueue_handler_with_options_records_the_author() {
2714 let mut engine = create_test_engine();
2715 engine.register(EchoWorkflow).unwrap();
2716 let actor = RunActor::User {
2717 user_id: Uuid::now_v7(),
2718 };
2719
2720 let run = engine
2721 .enqueue_handler_with_options(
2722 "echo-workflow",
2723 TriggerKind::Api,
2724 json!({}),
2725 EnqueueOptions {
2726 created_by: Some(actor.clone()),
2727 ..Default::default()
2728 },
2729 )
2730 .await
2731 .unwrap()
2732 .into_run();
2733
2734 assert_eq!(run.created_by, Some(actor));
2735 }
2736
2737 #[tokio::test]
2738 async fn enqueue_handler_with_options_accepts_no_author() {
2739 let mut engine = create_test_engine();
2740 engine.register(EchoWorkflow).unwrap();
2741
2742 let run = engine
2743 .enqueue_handler_with_options(
2744 "echo-workflow",
2745 TriggerKind::Cron {
2746 schedule: "0 * * * * *".to_string(),
2747 },
2748 json!({}),
2749 EnqueueOptions::default(),
2750 )
2751 .await
2752 .unwrap()
2753 .into_run();
2754
2755 assert!(run.created_by.is_none());
2756 }
2757
2758 #[tokio::test]
2759 async fn run_handler_leaves_the_run_unattributed() {
2760 let mut engine = create_test_engine();
2761 engine.register(EchoWorkflow).unwrap();
2762
2763 let run = engine
2764 .run_handler("echo-workflow", TriggerKind::Manual, json!({}))
2765 .await
2766 .unwrap()
2767 .run;
2768
2769 assert!(run.created_by.is_none());
2770 }
2771
2772 #[tokio::test]
2773 async fn engine_register_boxed() {
2774 let mut engine = create_test_engine();
2775 let handler: Box<dyn WorkflowHandler> = Box::new(EchoWorkflow);
2776 let result = engine.register_boxed(handler);
2777 assert!(result.is_ok());
2778 assert_eq!(engine.handler_names().len(), 1);
2779 }
2780
2781 #[tokio::test]
2782 async fn engine_store_and_provider_accessors() {
2783 let store = Arc::new(InMemoryStore::new());
2784 let inner = ClaudeCodeProvider::new();
2785 let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
2786 inner,
2787 "/tmp/ironflow-fixtures",
2788 ));
2789 let engine = Engine::new(store.clone(), provider.clone());
2790
2791 let _ = engine.store();
2793 let _ = engine.provider();
2794 }
2795
2796 use crate::operation::{Operation, OperationContext};
2801 use async_trait::async_trait;
2802 use ironflow_core::error::OperationError;
2803 use ironflow_store::models::StepKind;
2804
2805 struct FakeGitlabOp {
2806 project_id: u64,
2807 title: String,
2808 }
2809
2810 #[async_trait]
2811 impl Operation for FakeGitlabOp {
2812 fn kind(&self) -> &str {
2813 "gitlab"
2814 }
2815
2816 async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
2817 Ok(json!({
2818 "issue_id": 42,
2819 "project_id": self.project_id,
2820 "title": self.title,
2821 }))
2822 }
2823
2824 fn input(&self) -> Option<Value> {
2825 Some(json!({
2826 "project_id": self.project_id,
2827 "title": self.title,
2828 }))
2829 }
2830 }
2831
2832 struct FailingOp;
2833
2834 #[async_trait]
2835 impl Operation for FailingOp {
2836 fn kind(&self) -> &str {
2837 "broken-service"
2838 }
2839
2840 async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
2841 Err(OperationError::Http {
2842 status: None,
2843 message: "service unavailable".to_string(),
2844 })
2845 }
2846 }
2847
2848 struct OperationWorkflow;
2849
2850 impl WorkflowHandler for OperationWorkflow {
2851 fn name(&self) -> &str {
2852 "operation-workflow"
2853 }
2854
2855 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2856 Box::pin(async move {
2857 let op = FakeGitlabOp {
2858 project_id: 123,
2859 title: "Bug report".to_string(),
2860 };
2861 ctx.operation("create-issue", &op).await?;
2862 Ok(())
2863 })
2864 }
2865 }
2866
2867 struct FailingOperationWorkflow;
2868
2869 impl WorkflowHandler for FailingOperationWorkflow {
2870 fn name(&self) -> &str {
2871 "failing-operation-workflow"
2872 }
2873
2874 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2875 Box::pin(async move {
2876 ctx.operation("broken-call", &FailingOp).await?;
2877 Ok(())
2878 })
2879 }
2880 }
2881
2882 struct MixedWorkflow;
2883
2884 impl WorkflowHandler for MixedWorkflow {
2885 fn name(&self) -> &str {
2886 "mixed-workflow"
2887 }
2888
2889 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2890 Box::pin(async move {
2891 ctx.shell("build", ShellConfig::new("echo built")).await?;
2892 let op = FakeGitlabOp {
2893 project_id: 456,
2894 title: "Deploy done".to_string(),
2895 };
2896 let result = ctx.operation("notify-gitlab", &op).await?;
2897 assert_eq!(result.output["issue_id"], 42);
2898 Ok(())
2899 })
2900 }
2901 }
2902
2903 #[tokio::test]
2904 async fn operation_step_happy_path() {
2905 let mut engine = create_test_engine();
2906 engine.register(OperationWorkflow).unwrap();
2907
2908 let run = engine
2909 .run_handler("operation-workflow", TriggerKind::Manual, json!({}))
2910 .await
2911 .unwrap()
2912 .run;
2913
2914 assert_eq!(run.status.state, RunStatus::Completed);
2915
2916 let steps = engine.store().list_steps(run.id).await.unwrap();
2917
2918 assert_eq!(steps.len(), 1);
2919 assert_eq!(steps[0].name, "create-issue");
2920 assert_eq!(steps[0].kind, StepKind::Custom("gitlab".to_string()));
2921 assert_eq!(
2922 steps[0].status.state,
2923 ironflow_store::models::StepStatus::Completed
2924 );
2925
2926 let output = steps[0].output.as_ref().unwrap();
2927 assert_eq!(output["issue_id"], 42);
2928 assert_eq!(output["project_id"], 123);
2929
2930 let input = steps[0].input.as_ref().unwrap();
2931 assert_eq!(input["project_id"], 123);
2932 assert_eq!(input["title"], "Bug report");
2933 }
2934
2935 #[tokio::test]
2936 async fn operation_step_failure_marks_run_failed() {
2937 let mut engine = create_test_engine();
2938 engine.register(FailingOperationWorkflow).unwrap();
2939
2940 let result = engine
2941 .run_handler("failing-operation-workflow", TriggerKind::Manual, json!({}))
2942 .await;
2943
2944 assert!(result.is_err());
2945 }
2946
2947 #[tokio::test]
2948 async fn operation_mixed_with_shell_steps() {
2949 let mut engine = create_test_engine();
2950 engine.register(MixedWorkflow).unwrap();
2951
2952 let run = engine
2953 .run_handler("mixed-workflow", TriggerKind::Manual, json!({}))
2954 .await
2955 .unwrap()
2956 .run;
2957
2958 assert_eq!(run.status.state, RunStatus::Completed);
2959
2960 let steps = engine.store().list_steps(run.id).await.unwrap();
2961
2962 assert_eq!(steps.len(), 2);
2963 assert_eq!(steps[0].kind, StepKind::Shell);
2964 assert_eq!(steps[1].kind, StepKind::Custom("gitlab".to_string()));
2965 assert_eq!(steps[0].position, 0);
2966 assert_eq!(steps[1].position, 1);
2967 }
2968
2969 use crate::config::ApprovalConfig;
2974
2975 struct SingleApprovalWorkflow;
2976
2977 impl WorkflowHandler for SingleApprovalWorkflow {
2978 fn name(&self) -> &str {
2979 "single-approval"
2980 }
2981
2982 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2983 Box::pin(async move {
2984 ctx.shell("build", ShellConfig::new("echo built")).await?;
2985 ctx.approval("gate", ApprovalConfig::new("OK?")).await?;
2986 ctx.shell("deploy", ShellConfig::new("echo deployed"))
2987 .await?;
2988 Ok(())
2989 })
2990 }
2991 }
2992
2993 struct DoubleApprovalWorkflow;
2994
2995 impl WorkflowHandler for DoubleApprovalWorkflow {
2996 fn name(&self) -> &str {
2997 "double-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("staging-gate", ApprovalConfig::new("Deploy staging?"))
3004 .await?;
3005 ctx.shell("deploy-staging", ShellConfig::new("echo staging"))
3006 .await?;
3007 ctx.approval("prod-gate", ApprovalConfig::new("Deploy prod?"))
3008 .await?;
3009 ctx.shell("deploy-prod", ShellConfig::new("echo prod"))
3010 .await?;
3011 Ok(())
3012 })
3013 }
3014 }
3015
3016 #[tokio::test]
3017 async fn approval_pauses_run() {
3018 let mut engine = create_test_engine();
3019 engine.register(SingleApprovalWorkflow).unwrap();
3020
3021 let run = engine
3022 .run_handler("single-approval", TriggerKind::Manual, json!({}))
3023 .await
3024 .unwrap()
3025 .run;
3026
3027 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3028
3029 let steps = engine.store().list_steps(run.id).await.unwrap();
3030 assert_eq!(steps.len(), 2); assert_eq!(steps[0].kind, StepKind::Shell);
3032 assert_eq!(steps[0].status.state, StepStatus::Completed);
3033 assert_eq!(steps[1].kind, StepKind::Approval);
3034 assert_eq!(steps[1].status.state, StepStatus::AwaitingApproval);
3035 }
3036
3037 #[tokio::test]
3038 async fn approval_resume_completes_run() {
3039 let mut engine = create_test_engine();
3040 engine.register(SingleApprovalWorkflow).unwrap();
3041
3042 let run = engine
3044 .run_handler("single-approval", TriggerKind::Manual, json!({}))
3045 .await
3046 .unwrap()
3047 .run;
3048 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3049
3050 engine
3052 .store()
3053 .update_run_status(run.id, RunStatus::Running)
3054 .await
3055 .unwrap();
3056
3057 let resumed = engine.resume_run(run.id).await.unwrap().run;
3059 assert_eq!(resumed.status.state, RunStatus::Completed);
3060
3061 let steps = engine.store().list_steps(run.id).await.unwrap();
3062 assert_eq!(steps.len(), 3); assert_eq!(steps[0].name, "build");
3064 assert_eq!(steps[0].status.state, StepStatus::Completed);
3065 assert_eq!(steps[1].name, "gate");
3066 assert_eq!(steps[1].kind, StepKind::Approval);
3067 assert_eq!(steps[1].status.state, StepStatus::Completed);
3068 assert_eq!(steps[2].name, "deploy");
3069 assert_eq!(steps[2].status.state, StepStatus::Completed);
3070 }
3071
3072 #[tokio::test]
3073 async fn double_approval_two_resumes() {
3074 let mut engine = create_test_engine();
3075 engine.register(DoubleApprovalWorkflow).unwrap();
3076
3077 let run = engine
3079 .run_handler("double-approval", TriggerKind::Manual, json!({}))
3080 .await
3081 .unwrap()
3082 .run;
3083 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3084
3085 let steps = engine.store().list_steps(run.id).await.unwrap();
3086 assert_eq!(steps.len(), 2); engine
3090 .store()
3091 .update_run_status(run.id, RunStatus::Running)
3092 .await
3093 .unwrap();
3094
3095 let resumed = engine.resume_run(run.id).await.unwrap().run;
3096 assert_eq!(resumed.status.state, RunStatus::AwaitingApproval);
3097
3098 let steps = engine.store().list_steps(run.id).await.unwrap();
3099 assert_eq!(steps.len(), 4); engine
3103 .store()
3104 .update_run_status(run.id, RunStatus::Running)
3105 .await
3106 .unwrap();
3107
3108 let final_run = engine.resume_run(run.id).await.unwrap().run;
3109 assert_eq!(final_run.status.state, RunStatus::Completed);
3110
3111 let steps = engine.store().list_steps(run.id).await.unwrap();
3112 assert_eq!(steps.len(), 5);
3113 assert_eq!(steps[0].name, "build");
3114 assert_eq!(steps[1].name, "staging-gate");
3115 assert_eq!(steps[2].name, "deploy-staging");
3116 assert_eq!(steps[3].name, "prod-gate");
3117 assert_eq!(steps[4].name, "deploy-prod");
3118
3119 for step in &steps {
3120 assert_eq!(step.status.state, StepStatus::Completed);
3121 }
3122 }
3123
3124 use ironflow_store::models::{NewStep, StepUpdate, step_trace_id};
3129
3130 async fn create_step_with_status(
3131 store: &Arc<dyn Store>,
3132 run_id: Uuid,
3133 name: &str,
3134 position: u32,
3135 status: StepStatus,
3136 ) -> ironflow_store::models::Step {
3137 let step = store
3138 .create_step(NewStep {
3139 run_id,
3140 trace_id: step_trace_id(run_id, name, position),
3141 name: name.to_string(),
3142 kind: StepKind::Shell,
3143 position,
3144 input: None,
3145 is_error_handler: false,
3146 })
3147 .await
3148 .unwrap();
3149
3150 match status {
3151 StepStatus::Pending => {}
3152 StepStatus::Running => {
3153 store
3154 .update_step(
3155 step.id,
3156 StepUpdate {
3157 status: Some(StepStatus::Running),
3158 ..StepUpdate::default()
3159 },
3160 )
3161 .await
3162 .unwrap();
3163 }
3164 StepStatus::Completed => {
3165 store
3166 .update_step(
3167 step.id,
3168 StepUpdate {
3169 status: Some(StepStatus::Running),
3170 ..StepUpdate::default()
3171 },
3172 )
3173 .await
3174 .unwrap();
3175 store
3176 .update_step(
3177 step.id,
3178 StepUpdate {
3179 status: Some(StepStatus::Completed),
3180 ..StepUpdate::default()
3181 },
3182 )
3183 .await
3184 .unwrap();
3185 }
3186 StepStatus::AwaitingApproval => {
3187 store
3188 .update_step(
3189 step.id,
3190 StepUpdate {
3191 status: Some(StepStatus::Running),
3192 ..StepUpdate::default()
3193 },
3194 )
3195 .await
3196 .unwrap();
3197 store
3198 .update_step(
3199 step.id,
3200 StepUpdate {
3201 status: Some(StepStatus::AwaitingApproval),
3202 ..StepUpdate::default()
3203 },
3204 )
3205 .await
3206 .unwrap();
3207 }
3208 _ => panic!("unsupported status for test helper: {status}"),
3209 }
3210
3211 store.get_step(step.id).await.unwrap().unwrap()
3212 }
3213
3214 #[tokio::test]
3215 async fn fail_orphaned_steps_marks_running_as_failed() {
3216 let engine = create_test_engine();
3217 let run = engine
3218 .store()
3219 .create_run(NewRun {
3220 created_by: None,
3221 workflow_name: "test".to_string(),
3222 trigger: TriggerKind::Manual,
3223 payload: json!({}),
3224 max_retries: 0,
3225 handler_version: None,
3226 labels: HashMap::new(),
3227 scheduled_at: None,
3228 idempotency_key: None,
3229 concurrency_key: None,
3230 max_cost_usd: None,
3231 })
3232 .await
3233 .unwrap()
3234 .into_run();
3235
3236 let step = create_step_with_status(
3237 engine.store(),
3238 run.id,
3239 "running-step",
3240 0,
3241 StepStatus::Running,
3242 )
3243 .await;
3244
3245 engine
3246 .fail_orphaned_steps(run.id, "parent run timed out")
3247 .await
3248 .unwrap();
3249
3250 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
3251 assert_eq!(updated.status.state, StepStatus::Failed);
3252 assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
3253 assert!(updated.completed_at.is_some());
3254 }
3255
3256 #[tokio::test]
3257 async fn fail_orphaned_steps_marks_pending_as_skipped() {
3258 let engine = create_test_engine();
3259 let run = engine
3260 .store()
3261 .create_run(NewRun {
3262 created_by: None,
3263 workflow_name: "test".to_string(),
3264 trigger: TriggerKind::Manual,
3265 payload: json!({}),
3266 max_retries: 0,
3267 handler_version: None,
3268 labels: HashMap::new(),
3269 scheduled_at: None,
3270 idempotency_key: None,
3271 concurrency_key: None,
3272 max_cost_usd: None,
3273 })
3274 .await
3275 .unwrap()
3276 .into_run();
3277
3278 let step = create_step_with_status(
3279 engine.store(),
3280 run.id,
3281 "pending-step",
3282 0,
3283 StepStatus::Pending,
3284 )
3285 .await;
3286
3287 engine
3288 .fail_orphaned_steps(run.id, "parent run timed out")
3289 .await
3290 .unwrap();
3291
3292 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
3293 assert_eq!(updated.status.state, StepStatus::Skipped);
3294 assert!(updated.error.is_none());
3295 assert!(updated.completed_at.is_some());
3296 }
3297
3298 #[tokio::test]
3299 async fn fail_orphaned_steps_marks_awaiting_approval_as_failed() {
3300 let engine = create_test_engine();
3301 let run = engine
3302 .store()
3303 .create_run(NewRun {
3304 created_by: None,
3305 workflow_name: "test".to_string(),
3306 trigger: TriggerKind::Manual,
3307 payload: json!({}),
3308 max_retries: 0,
3309 handler_version: None,
3310 labels: HashMap::new(),
3311 scheduled_at: None,
3312 idempotency_key: None,
3313 concurrency_key: None,
3314 max_cost_usd: None,
3315 })
3316 .await
3317 .unwrap()
3318 .into_run();
3319
3320 let step = create_step_with_status(
3321 engine.store(),
3322 run.id,
3323 "approval-step",
3324 0,
3325 StepStatus::AwaitingApproval,
3326 )
3327 .await;
3328
3329 engine
3330 .fail_orphaned_steps(run.id, "parent run timed out")
3331 .await
3332 .unwrap();
3333
3334 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
3335 assert_eq!(updated.status.state, StepStatus::Failed);
3336 assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
3337 assert!(updated.completed_at.is_some());
3338 }
3339
3340 #[tokio::test]
3341 async fn fail_orphaned_steps_skips_terminal_steps() {
3342 let engine = create_test_engine();
3343 let run = engine
3344 .store()
3345 .create_run(NewRun {
3346 created_by: None,
3347 workflow_name: "test".to_string(),
3348 trigger: TriggerKind::Manual,
3349 payload: json!({}),
3350 max_retries: 0,
3351 handler_version: None,
3352 labels: HashMap::new(),
3353 scheduled_at: None,
3354 idempotency_key: None,
3355 concurrency_key: None,
3356 max_cost_usd: None,
3357 })
3358 .await
3359 .unwrap()
3360 .into_run();
3361
3362 let completed_step =
3363 create_step_with_status(engine.store(), run.id, "done", 0, StepStatus::Completed).await;
3364 let running_step =
3365 create_step_with_status(engine.store(), run.id, "in-flight", 1, StepStatus::Running)
3366 .await;
3367
3368 engine
3369 .fail_orphaned_steps(run.id, "parent run timed out")
3370 .await
3371 .unwrap();
3372
3373 let completed = engine
3374 .store()
3375 .get_step(completed_step.id)
3376 .await
3377 .unwrap()
3378 .unwrap();
3379 assert_eq!(completed.status.state, StepStatus::Completed);
3380
3381 let failed = engine
3382 .store()
3383 .get_step(running_step.id)
3384 .await
3385 .unwrap()
3386 .unwrap();
3387 assert_eq!(failed.status.state, StepStatus::Failed);
3388 }
3389
3390 #[tokio::test]
3391 async fn fail_orphaned_steps_mixed_states() {
3392 let engine = create_test_engine();
3393 let run = engine
3394 .store()
3395 .create_run(NewRun {
3396 created_by: None,
3397 workflow_name: "test".to_string(),
3398 trigger: TriggerKind::Manual,
3399 payload: json!({}),
3400 max_retries: 0,
3401 handler_version: None,
3402 labels: HashMap::new(),
3403 scheduled_at: None,
3404 idempotency_key: None,
3405 concurrency_key: None,
3406 max_cost_usd: None,
3407 })
3408 .await
3409 .unwrap()
3410 .into_run();
3411
3412 let s_completed =
3413 create_step_with_status(engine.store(), run.id, "step-1", 0, StepStatus::Completed)
3414 .await;
3415 let s_running =
3416 create_step_with_status(engine.store(), run.id, "step-2", 1, StepStatus::Running).await;
3417 let s_pending =
3418 create_step_with_status(engine.store(), run.id, "step-3", 2, StepStatus::Pending).await;
3419
3420 engine.fail_orphaned_steps(run.id, "timeout").await.unwrap();
3421
3422 let r_completed = engine
3423 .store()
3424 .get_step(s_completed.id)
3425 .await
3426 .unwrap()
3427 .unwrap();
3428 assert_eq!(r_completed.status.state, StepStatus::Completed);
3429
3430 let r_running = engine
3431 .store()
3432 .get_step(s_running.id)
3433 .await
3434 .unwrap()
3435 .unwrap();
3436 assert_eq!(r_running.status.state, StepStatus::Failed);
3437 assert_eq!(r_running.error.as_deref(), Some("timeout"));
3438
3439 let r_pending = engine
3440 .store()
3441 .get_step(s_pending.id)
3442 .await
3443 .unwrap()
3444 .unwrap();
3445 assert_eq!(r_pending.status.state, StepStatus::Skipped);
3446 assert!(r_pending.error.is_none());
3447 }
3448
3449 #[tokio::test]
3450 async fn fail_orphaned_steps_no_steps_is_noop() {
3451 let engine = create_test_engine();
3452 let run = engine
3453 .store()
3454 .create_run(NewRun {
3455 created_by: None,
3456 workflow_name: "test".to_string(),
3457 trigger: TriggerKind::Manual,
3458 payload: json!({}),
3459 max_retries: 0,
3460 handler_version: None,
3461 labels: HashMap::new(),
3462 scheduled_at: None,
3463 idempotency_key: None,
3464 concurrency_key: None,
3465 max_cost_usd: None,
3466 })
3467 .await
3468 .unwrap()
3469 .into_run();
3470
3471 let result = engine.fail_orphaned_steps(run.id, "timeout").await;
3472 assert!(result.is_ok());
3473 }
3474
3475 #[tokio::test]
3476 async fn fail_orphaned_steps_preserves_existing_error() {
3477 let engine = create_test_engine();
3478 let run = engine
3479 .store()
3480 .create_run(NewRun {
3481 created_by: None,
3482 workflow_name: "test".to_string(),
3483 trigger: TriggerKind::Manual,
3484 payload: json!({}),
3485 max_retries: 0,
3486 handler_version: None,
3487 labels: HashMap::new(),
3488 scheduled_at: None,
3489 idempotency_key: None,
3490 concurrency_key: None,
3491 max_cost_usd: None,
3492 })
3493 .await
3494 .unwrap()
3495 .into_run();
3496
3497 let step_with_error = create_step_with_status(
3498 engine.store(),
3499 run.id,
3500 "already-errored",
3501 0,
3502 StepStatus::Running,
3503 )
3504 .await;
3505
3506 engine
3507 .store()
3508 .update_step(
3509 step_with_error.id,
3510 StepUpdate {
3511 error: Some("real error from provider".to_string()),
3512 ..StepUpdate::default()
3513 },
3514 )
3515 .await
3516 .unwrap();
3517
3518 let step_no_error = create_step_with_status(
3519 engine.store(),
3520 run.id,
3521 "no-error-yet",
3522 1,
3523 StepStatus::Running,
3524 )
3525 .await;
3526
3527 engine
3528 .fail_orphaned_steps(run.id, "parent run failed")
3529 .await
3530 .unwrap();
3531
3532 let updated_with = engine
3533 .store()
3534 .get_step(step_with_error.id)
3535 .await
3536 .unwrap()
3537 .unwrap();
3538 assert_eq!(updated_with.status.state, StepStatus::Failed);
3539 assert_eq!(
3540 updated_with.error.as_deref(),
3541 Some("real error from provider"),
3542 );
3543
3544 let updated_without = engine
3545 .store()
3546 .get_step(step_no_error.id)
3547 .await
3548 .unwrap()
3549 .unwrap();
3550 assert_eq!(updated_without.status.state, StepStatus::Failed);
3551 assert_eq!(updated_without.error.as_deref(), Some("parent run failed"),);
3552 }
3553}