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