1use std::collections::HashMap;
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;
17use tracing::{error, info, warn};
18use uuid::Uuid;
19
20use ironflow_core::error::OperationError;
21#[cfg(feature = "prometheus")]
22use ironflow_core::metric_names::{
23 RUN_BUDGET_EXCEEDED_TOTAL, RUN_COST_USD, RUN_DURATION_SECONDS, RUNS_ACTIVE, RUNS_TOTAL,
24};
25use ironflow_core::provider::AgentProvider;
26use ironflow_store::error::StoreError;
27use ironflow_store::models::{
28 NewRun, Run, RunActor, RunCreation, RunFilter, RunStatus, RunUpdate, StepStatus, StepUpdate,
29 TriggerKind,
30};
31use ironflow_store::store::Store;
32#[cfg(feature = "prometheus")]
33use metrics::{counter, gauge, histogram};
34
35use crate::artifact::ArtifactSink;
36use crate::budget::{BudgetConfig, month_start};
37use crate::context::WorkflowContext;
38use crate::error::EngineError;
39use crate::executor::{StepInterceptor, StepResult};
40use crate::guard::{WorkflowGuardConfig, new_shared_guard_state};
41use crate::handler::{WorkflowHandler, WorkflowInfo};
42use crate::log_sender::LogSender;
43use crate::notify::{
44 ApprovalRequestedEvent, Event, EventPublisher, EventSubscriber, RunBudgetExceededEvent,
45 RunFailedEvent, RunStatusChangedEvent, WorkflowEventBus,
46};
47use crate::plan::{
48 ExecutionPlan, PlanOptions, PlanRecorder, SharedPlanRecorder, estimate_durations, lock_plan,
49};
50use crate::retry_policy::{backoff_for_retry, is_run_retryable};
51use crate::schedule::CronSchedule;
52use ironflow_core::decision::DecisionProvider;
53
54#[derive(Debug, Clone)]
73pub struct WorkflowResult {
74 pub run: Run,
76 pub steps: Vec<StepResult>,
78}
79
80#[derive(Debug, Clone, Default)]
99pub struct EnqueueOptions {
100 pub max_retries: u32,
102 pub labels: HashMap<String, String>,
104 pub scheduled_at: Option<DateTime<Utc>>,
107 pub max_cost_usd: Option<Decimal>,
111 pub created_by: Option<RunActor>,
114 pub idempotency_key: Option<String>,
120}
121
122#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
133pub enum ExecutionMode {
134 #[default]
138 Local,
139 Workers,
142}
143
144pub struct Engine {
185 store: Arc<dyn Store>,
186 provider: Arc<dyn AgentProvider>,
187 handlers: HashMap<String, Arc<dyn WorkflowHandler>>,
188 event_publisher: EventPublisher,
189 log_sender: Option<LogSender>,
190 budget: BudgetConfig,
191 artifact_sink: Option<Arc<dyn ArtifactSink>>,
192 guard_config: Option<WorkflowGuardConfig>,
193 event_bus: Option<WorkflowEventBus>,
194 decision_provider: Option<Arc<dyn DecisionProvider>>,
195 step_interceptor: Option<Arc<dyn StepInterceptor>>,
196 execution_mode: ExecutionMode,
197}
198
199fn validate_category(handler_name: &str, category: &str) -> Result<(), EngineError> {
209 let reject = |reason: &str| {
210 Err(EngineError::InvalidWorkflow(format!(
211 "handler '{handler_name}' has invalid category '{category}': {reason}"
212 )))
213 };
214
215 if category.is_empty() {
216 return reject("empty category");
217 }
218 if category.starts_with('/') {
219 return reject("leading '/'");
220 }
221 if category.ends_with('/') {
222 return reject("trailing '/'");
223 }
224 for segment in category.split('/') {
225 if segment.is_empty() {
226 return reject("empty segment (double '/')");
227 }
228 if segment.trim().is_empty() {
229 return reject("whitespace-only segment");
230 }
231 }
232 Ok(())
233}
234
235impl Engine {
236 pub fn new(store: Arc<dyn Store>, provider: Arc<dyn AgentProvider>) -> Self {
252 Self {
253 store,
254 provider,
255 handlers: HashMap::new(),
256 event_publisher: EventPublisher::new(),
257 log_sender: None,
258 budget: BudgetConfig::new(),
259 artifact_sink: None,
260 guard_config: None,
261 event_bus: None,
262 decision_provider: None,
263 step_interceptor: None,
264 execution_mode: ExecutionMode::default(),
265 }
266 }
267
268 pub fn with_decision_provider(mut self, provider: Arc<dyn DecisionProvider>) -> Self {
289 self.decision_provider = Some(provider);
290 self
291 }
292
293 pub fn with_step_interceptor(mut self, interceptor: Arc<dyn StepInterceptor>) -> Self {
317 self.step_interceptor = Some(interceptor);
318 self
319 }
320
321 pub fn step_interceptor(&self) -> Option<&Arc<dyn StepInterceptor>> {
323 self.step_interceptor.as_ref()
324 }
325
326 pub fn with_budget_config(mut self, budget: BudgetConfig) -> Self {
347 self.budget = budget;
348 self
349 }
350
351 pub fn budget_config(&self) -> &BudgetConfig {
353 &self.budget
354 }
355
356 pub fn with_guard_config(mut self, config: WorkflowGuardConfig) -> Self {
378 self.guard_config = Some(config);
379 self
380 }
381
382 pub fn guard_config(&self) -> Option<&WorkflowGuardConfig> {
384 self.guard_config.as_ref()
385 }
386
387 pub fn with_execution_mode(mut self, mode: ExecutionMode) -> Self {
408 self.execution_mode = mode;
409 self
410 }
411
412 pub fn execution_mode(&self) -> ExecutionMode {
414 self.execution_mode
415 }
416
417 pub fn set_log_sender(&mut self, sender: LogSender) {
423 self.log_sender = Some(sender);
424 }
425
426 pub fn set_artifact_sink(&mut self, sink: Arc<dyn ArtifactSink>) {
444 self.artifact_sink = Some(sink);
445 }
446
447 pub fn artifact_sink(&self) -> Option<&Arc<dyn ArtifactSink>> {
449 self.artifact_sink.as_ref()
450 }
451
452 pub fn set_event_bus(&mut self, bus: WorkflowEventBus) {
469 self.event_bus = Some(bus);
470 }
471
472 pub fn event_bus(&self) -> Option<&WorkflowEventBus> {
474 self.event_bus.as_ref()
475 }
476
477 pub fn store(&self) -> &Arc<dyn Store> {
479 &self.store
480 }
481
482 pub fn provider(&self) -> &Arc<dyn AgentProvider> {
484 &self.provider
485 }
486
487 fn build_context(&self, run: &Run) -> WorkflowContext {
496 let handlers = self.handlers.clone();
497 let resolver: crate::context::HandlerResolver =
498 Arc::new(move |name: &str| handlers.get(name).cloned());
499 let mut ctx = WorkflowContext::with_handler_resolver(
500 run.id,
501 run.workflow_name.clone(),
502 self.store.clone(),
503 self.provider.clone(),
504 resolver,
505 );
506 ctx.carry_over_run_totals(run.retry_count + 1, run.cost_usd, run.duration_ms);
507 ctx.set_max_cost_usd(run.max_cost_usd);
508 if let Some(ref sender) = self.log_sender {
509 ctx.set_log_sender(sender.clone());
510 }
511 if let Some(ref sink) = self.artifact_sink {
512 ctx.set_artifact_sink(sink.clone());
513 }
514 if let Some(ref bus) = self.event_bus {
515 ctx.set_event_bus(bus.clone());
516 }
517 if let Some(ref provider) = self.decision_provider {
518 ctx.set_decision_provider(provider.clone());
519 }
520 if let Some(ref interceptor) = self.step_interceptor {
521 ctx.set_step_interceptor(interceptor.clone());
522 }
523 ctx
524 }
525
526 fn build_context_with_guard(
532 &self,
533 run: &Run,
534 handler: &dyn WorkflowHandler,
535 ) -> WorkflowContext {
536 let mut ctx = self.build_context(run);
537 let guard_config = handler.guard_config().or_else(|| self.guard_config.clone());
538 if let Some(config) = guard_config {
539 ctx.set_guard(config, new_shared_guard_state());
540 }
541 ctx
542 }
543
544 async fn check_monthly_quota(&self, workflow_name: &str) -> Result<(), EngineError> {
555 let Some(limit) = self.budget.monthly_cost_limit_usd else {
556 return Ok(());
557 };
558
559 let stats = self
560 .store
561 .get_stats(RunFilter {
562 created_after: Some(month_start(Utc::now())),
563 ..RunFilter::default()
564 })
565 .await?;
566
567 if stats.total_cost_usd < limit {
568 return Ok(());
569 }
570
571 warn!(
572 workflow = %workflow_name,
573 limit_usd = %limit,
574 spent_usd = %stats.total_cost_usd,
575 "monthly cost quota exhausted, refusing new run"
576 );
577
578 #[cfg(feature = "prometheus")]
579 counter!(
580 RUN_BUDGET_EXCEEDED_TOTAL,
581 "workflow" => workflow_name.to_string(),
582 "scope" => "monthly",
583 )
584 .increment(1);
585
586 Err(EngineError::MonthlyBudgetExceeded {
587 limit_usd: limit,
588 spent_usd: stats.total_cost_usd,
589 })
590 }
591
592 pub fn register(&mut self, handler: impl WorkflowHandler + 'static) -> Result<(), EngineError> {
636 let name = handler.name().to_string();
637 if self.handlers.contains_key(&name) {
638 return Err(EngineError::InvalidWorkflow(format!(
639 "handler '{}' already registered",
640 name
641 )));
642 }
643 if let Some(category) = handler.category() {
644 validate_category(&name, category)?;
645 }
646 self.handlers.insert(name, Arc::new(handler));
647 Ok(())
648 }
649
650 pub fn register_boxed(&mut self, handler: Box<dyn WorkflowHandler>) -> Result<(), EngineError> {
657 let name = handler.name().to_string();
658 if self.handlers.contains_key(&name) {
659 return Err(EngineError::InvalidWorkflow(format!(
660 "handler '{}' already registered",
661 name
662 )));
663 }
664 if let Some(category) = handler.category() {
665 validate_category(&name, category)?;
666 }
667 self.handlers.insert(name, Arc::from(handler));
668 Ok(())
669 }
670
671 pub fn get_handler(&self, name: &str) -> Option<&Arc<dyn WorkflowHandler>> {
673 self.handlers.get(name)
674 }
675
676 pub fn handler_names(&self) -> Vec<&str> {
678 self.handlers.keys().map(|s| s.as_str()).collect()
679 }
680
681 pub fn handler_info(&self, name: &str) -> Option<WorkflowInfo> {
683 self.handlers.get(name).map(|h| h.describe())
684 }
685
686 pub fn scheduled_handlers(&self) -> Vec<(&str, &CronSchedule)> {
710 self.handlers
711 .iter()
712 .filter_map(|(name, handler)| handler.schedule().map(|sched| (name.as_str(), sched)))
713 .collect()
714 }
715
716 pub fn subscribe(
741 &mut self,
742 subscriber: impl EventSubscriber + 'static,
743 event_types: &[&'static str],
744 ) {
745 self.event_publisher.subscribe(subscriber, event_types);
746 }
747
748 pub fn event_publisher(&self) -> &EventPublisher {
753 &self.event_publisher
754 }
755
756 #[tracing::instrument(name = "engine.run_handler", skip_all, fields(workflow = %handler_name))]
786 pub async fn run_handler(
787 &self,
788 handler_name: &str,
789 trigger: TriggerKind,
790 payload: Value,
791 ) -> Result<WorkflowResult, EngineError> {
792 let handler = self
793 .handlers
794 .get(handler_name)
795 .ok_or_else(|| {
796 EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
797 })?
798 .clone();
799
800 self.check_monthly_quota(handler_name).await?;
801
802 let handler_version = handler.version().map(str::to_string);
803 let max_cost_usd = self
804 .budget
805 .resolve_run_cap(None, handler.default_max_cost_usd());
806 let run = self
807 .store
808 .create_run(NewRun {
809 created_by: None,
810 workflow_name: handler_name.to_string(),
811 trigger,
812 payload,
813 max_retries: 0,
814 handler_version,
815 labels: handler.default_labels(),
816 scheduled_at: None,
817 idempotency_key: None,
818 max_cost_usd,
819 })
820 .await?
821 .into_run();
822
823 let run_id = run.id;
824 info!(run_id = %run_id, handler_version = run.handler_version.as_deref().unwrap_or(""), "run created");
825
826 self.store
827 .update_run_status(run_id, RunStatus::Running)
828 .await?;
829
830 #[cfg(feature = "prometheus")]
831 gauge!(RUNS_ACTIVE, "workflow" => handler_name.to_string()).increment(1.0);
832
833 let run_start = Instant::now();
834 let mut ctx = self.build_context_with_guard(&run, handler.as_ref());
835
836 let result = handler.execute(&mut ctx).await;
837 self.finalize_run(run_id, handler_name, result, &ctx, run_start, run.labels)
838 .await
839 }
840
841 #[tracing::instrument(name = "engine.plan_handler", skip_all, fields(workflow = %handler_name))]
883 pub async fn plan_handler(
884 &self,
885 handler_name: &str,
886 payload: Value,
887 options: PlanOptions,
888 ) -> Result<ExecutionPlan, EngineError> {
889 if options.max_depth == 0 {
890 return Err(EngineError::InvalidWorkflow(
891 "max_depth must be at least 1".to_string(),
892 ));
893 }
894
895 let handler = self
896 .handlers
897 .get(handler_name)
898 .ok_or_else(|| {
899 EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
900 })?
901 .clone();
902
903 let estimates = if options.estimate_durations {
904 estimate_durations(&self.store, handler_name, options.sample_runs).await?
905 } else {
906 HashMap::new()
907 };
908
909 let shared: SharedPlanRecorder = Arc::new(Mutex::new(PlanRecorder::new(
910 handler_name.to_string(),
911 payload,
912 options.max_depth,
913 estimates,
914 )));
915
916 let handlers = self.handlers.clone();
919 let resolver: crate::context::HandlerResolver =
920 Arc::new(move |name: &str| handlers.get(name).cloned());
921 let mut ctx = WorkflowContext::with_handler_resolver(
922 Uuid::now_v7(),
923 handler_name.to_string(),
924 self.store.clone(),
925 self.provider.clone(),
926 resolver,
927 );
928 ctx.set_plan(shared.clone());
929
930 if let Err(err) = handler.execute(&mut ctx).await {
931 lock_plan(&shared).fail(err.to_string());
932 }
933 drop(ctx);
934
935 let plan = match Arc::try_unwrap(shared) {
936 Ok(mutex) => mutex
937 .into_inner()
938 .unwrap_or_else(|poisoned| poisoned.into_inner())
939 .into_plan(),
940 Err(shared) => lock_plan(&shared).snapshot(),
941 };
942
943 info!(
944 workflow = %handler_name,
945 steps = plan.steps.len(),
946 truncated = plan.truncated,
947 "execution plan built"
948 );
949
950 Ok(plan)
951 }
952
953 #[tracing::instrument(name = "engine.enqueue_handler", skip_all, fields(workflow = %handler_name))]
964 pub async fn enqueue_handler(
965 &self,
966 handler_name: &str,
967 trigger: TriggerKind,
968 payload: Value,
969 max_retries: u32,
970 ) -> Result<Run, EngineError> {
971 self.enqueue_handler_with_options(
972 handler_name,
973 trigger,
974 payload,
975 EnqueueOptions {
976 max_retries,
977 ..Default::default()
978 },
979 )
980 .await
981 .map(RunCreation::into_run)
982 }
983
984 #[tracing::instrument(name = "engine.enqueue_handler_with_options", skip_all, fields(workflow = %handler_name))]
1029 pub async fn enqueue_handler_with_options(
1030 &self,
1031 handler_name: &str,
1032 trigger: TriggerKind,
1033 payload: Value,
1034 options: EnqueueOptions,
1035 ) -> Result<RunCreation, EngineError> {
1036 let EnqueueOptions {
1037 max_retries,
1038 labels,
1039 scheduled_at,
1040 max_cost_usd,
1041 created_by,
1042 idempotency_key,
1043 } = options;
1044
1045 let handler = self.handlers.get(handler_name).ok_or_else(|| {
1046 EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
1047 })?;
1048
1049 self.check_monthly_quota(handler_name).await?;
1050
1051 let handler_version = handler.version().map(str::to_string);
1052 let mut merged_labels = handler.default_labels();
1053 merged_labels.extend(labels);
1054 let resolved_cap = self
1055 .budget
1056 .resolve_run_cap(max_cost_usd, handler.default_max_cost_usd());
1057
1058 let creation = self
1059 .store
1060 .create_run(NewRun {
1061 workflow_name: handler_name.to_string(),
1062 trigger,
1063 payload,
1064 max_retries,
1065 handler_version,
1066 labels: merged_labels,
1067 scheduled_at,
1068 created_by,
1069 idempotency_key,
1070 max_cost_usd: resolved_cap,
1071 })
1072 .await?;
1073
1074 match &creation {
1075 RunCreation::Created(run) => info!(
1076 run_id = %run.id,
1077 workflow = %handler_name,
1078 max_cost_usd = ?resolved_cap,
1079 "handler run enqueued"
1080 ),
1081 RunCreation::Existing(run) => info!(
1082 run_id = %run.id,
1083 workflow = %handler_name,
1084 "idempotent replay, nothing enqueued"
1085 ),
1086 }
1087
1088 Ok(creation)
1089 }
1090
1091 #[tracing::instrument(name = "engine.execute_handler_run", skip_all, fields(run_id = %run_id))]
1107 pub async fn execute_handler_run(&self, run_id: Uuid) -> Result<WorkflowResult, EngineError> {
1108 let run = self
1109 .store
1110 .get_run(run_id)
1111 .await?
1112 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
1113
1114 let handler = self
1115 .handlers
1116 .get(&run.workflow_name)
1117 .ok_or_else(|| {
1118 EngineError::InvalidWorkflow(format!(
1119 "no handler registered: {}",
1120 run.workflow_name
1121 ))
1122 })?
1123 .clone();
1124
1125 #[cfg(feature = "prometheus")]
1126 gauge!(RUNS_ACTIVE, "workflow" => run.workflow_name.clone()).increment(1.0);
1127
1128 let run_start = Instant::now();
1129 let mut ctx = self.build_context_with_guard(&run, handler.as_ref());
1130
1131 let result = if handler.is_version_compatible(run.handler_version.as_deref()) {
1142 ctx.load_replay_steps().await?;
1143 self.release_then_execute(run_id, handler.as_ref(), &mut ctx)
1144 .await
1145 } else {
1146 Err(EngineError::HandlerVersionMismatch {
1147 run_id,
1148 workflow_name: run.workflow_name.clone(),
1149 run_version: run
1150 .handler_version
1151 .clone()
1152 .unwrap_or_else(|| "unknown".to_string()),
1153 current_version: handler
1154 .version()
1155 .map(str::to_string)
1156 .unwrap_or_else(|| "unknown".to_string()),
1157 })
1158 };
1159
1160 self.finalize_run(
1161 run_id,
1162 &run.workflow_name,
1163 result,
1164 &ctx,
1165 run_start,
1166 run.labels,
1167 )
1168 .await
1169 }
1170
1171 #[tracing::instrument(name = "engine.execute_run", skip_all, fields(run_id = %run_id))]
1179 pub async fn execute_run(&self, run_id: Uuid) -> Result<WorkflowResult, EngineError> {
1180 self.execute_handler_run(run_id).await
1181 }
1182
1183 #[tracing::instrument(name = "engine.resume_run", skip_all, fields(run_id = %run_id))]
1205 pub async fn resume_run(&self, run_id: Uuid) -> Result<WorkflowResult, EngineError> {
1206 let run = self
1207 .store
1208 .get_run(run_id)
1209 .await?
1210 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
1211
1212 let handler = self
1213 .handlers
1214 .get(&run.workflow_name)
1215 .ok_or_else(|| {
1216 EngineError::InvalidWorkflow(format!(
1217 "no handler registered: {}",
1218 run.workflow_name
1219 ))
1220 })?
1221 .clone();
1222
1223 info!(run_id = %run_id, workflow = %run.workflow_name, "resuming run after approval");
1224
1225 let run_start = Instant::now();
1226 let mut ctx = self.build_context_with_guard(&run, handler.as_ref());
1227
1228 let result = if handler.is_version_compatible(run.handler_version.as_deref()) {
1229 ctx.load_replay_steps().await?;
1230 self.release_then_execute(run_id, handler.as_ref(), &mut ctx)
1231 .await
1232 } else {
1233 Err(EngineError::HandlerVersionMismatch {
1234 run_id,
1235 workflow_name: run.workflow_name.clone(),
1236 run_version: run
1237 .handler_version
1238 .clone()
1239 .unwrap_or_else(|| "unknown".to_string()),
1240 current_version: handler
1241 .version()
1242 .map(str::to_string)
1243 .unwrap_or_else(|| "unknown".to_string()),
1244 })
1245 };
1246
1247 self.finalize_run(
1248 run_id,
1249 &run.workflow_name,
1250 result,
1251 &ctx,
1252 run_start,
1253 run.labels,
1254 )
1255 .await
1256 }
1257
1258 pub async fn fail_or_schedule_retry(
1303 &self,
1304 run_id: Uuid,
1305 error: &str,
1306 retryable: bool,
1307 cost_usd: Option<Decimal>,
1308 duration_ms: Option<u64>,
1309 ) -> Result<RunStatus, EngineError> {
1310 let run = self
1311 .store
1312 .get_run(run_id)
1313 .await?
1314 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
1315
1316 let has_attempts_left = run.retry_count < run.max_retries;
1317 let update = if retryable && has_attempts_left {
1318 let backoff = backoff_for_retry(run.retry_count);
1319 let scheduled_at = Utc::now() + TimeDelta::milliseconds(backoff.as_millis() as i64);
1320
1321 info!(
1322 run_id = %run_id,
1323 workflow = %run.workflow_name,
1324 attempt = run.retry_count + 1,
1325 max_retries = run.max_retries,
1326 backoff_secs = backoff.as_secs(),
1327 scheduled_at = %scheduled_at,
1328 "run failed, scheduling retry"
1329 );
1330
1331 RunUpdate {
1332 status: Some(RunStatus::Retrying),
1333 error: Some(error.to_string()),
1334 increment_retry: true,
1335 cost_usd,
1336 duration_ms,
1337 scheduled_at: Some(scheduled_at),
1338 ..RunUpdate::default()
1339 }
1340 } else {
1341 RunUpdate {
1342 status: Some(RunStatus::Failed),
1343 error: Some(error.to_string()),
1344 cost_usd,
1345 duration_ms,
1346 completed_at: Some(Utc::now()),
1347 ..RunUpdate::default()
1348 }
1349 };
1350
1351 let status = update.status.unwrap_or(RunStatus::Failed);
1352 self.store.update_run(run_id, update).await?;
1353 self.fail_orphaned_steps(run_id, error).await?;
1354
1355 Ok(status)
1356 }
1357
1358 pub async fn fail_orphaned_steps(
1372 &self,
1373 run_id: Uuid,
1374 error_message: &str,
1375 ) -> Result<(), EngineError> {
1376 let steps = self.store.list_steps(run_id).await?;
1377 let now = Utc::now();
1378
1379 for step in steps {
1380 if step.status.state.is_terminal() {
1381 continue;
1382 }
1383
1384 let (target_status, error) = match step.status.state {
1385 StepStatus::Running | StepStatus::AwaitingApproval => {
1386 let err = if step.error.is_some() {
1387 None
1388 } else {
1389 Some(error_message.to_string())
1390 };
1391 (StepStatus::Failed, err)
1392 }
1393 StepStatus::Pending => (StepStatus::Skipped, None),
1394 _ => continue,
1395 };
1396
1397 if let Err(e) = self
1398 .store
1399 .update_step(
1400 step.id,
1401 StepUpdate {
1402 status: Some(target_status),
1403 error,
1404 completed_at: Some(now),
1405 ..StepUpdate::default()
1406 },
1407 )
1408 .await
1409 {
1410 warn!(
1411 run_id = %run_id,
1412 step_id = %step.id,
1413 step_name = %step.name,
1414 error = %e,
1415 "failed to cleanup orphaned step"
1416 );
1417 } else {
1418 info!(
1419 run_id = %run_id,
1420 step_id = %step.id,
1421 step_name = %step.name,
1422 from = %step.status.state,
1423 to = %target_status,
1424 "cleaned up orphaned step"
1425 );
1426 }
1427 }
1428
1429 Ok(())
1430 }
1431
1432 async fn release_then_execute(
1437 &self,
1438 run_id: Uuid,
1439 handler: &dyn WorkflowHandler,
1440 ctx: &mut WorkflowContext,
1441 ) -> Result<(), EngineError> {
1442 match self.provider.release_run(&run_id.to_string()).await {
1443 Ok(()) => handler.execute(ctx).await,
1444 Err(e) => Err(EngineError::Operation(OperationError::Agent(e))),
1445 }
1446 }
1447
1448 async fn finalize_run(
1454 &self,
1455 run_id: Uuid,
1456 workflow_name: &str,
1457 result: Result<(), EngineError>,
1458 ctx: &WorkflowContext,
1459 run_start: Instant,
1460 run_labels: HashMap<String, String>,
1461 ) -> Result<WorkflowResult, EngineError> {
1462 let total_duration = ctx.carried_duration_ms() + run_start.elapsed().as_millis() as u64;
1465 let completed_at = Utc::now();
1466
1467 let final_status;
1468 let final_run;
1469
1470 match result {
1471 Ok(()) => {
1472 final_status = if ctx.has_allowed_failure() {
1473 RunStatus::Warning
1474 } else {
1475 RunStatus::Completed
1476 };
1477 final_run = self
1478 .store
1479 .update_run_returning(
1480 run_id,
1481 RunUpdate {
1482 status: Some(final_status),
1483 cost_usd: Some(ctx.total_cost_usd()),
1484 duration_ms: Some(total_duration),
1485 completed_at: Some(completed_at),
1486 ..RunUpdate::default()
1487 },
1488 )
1489 .await?;
1490
1491 info!(
1492 run_id = %run_id,
1493 status = %final_status,
1494 cost_usd = %ctx.total_cost_usd(),
1495 duration_ms = total_duration,
1496 "run completed"
1497 );
1498 }
1499 Err(EngineError::ApprovalRequired {
1500 run_id: approval_run_id,
1501 step_id,
1502 ref message,
1503 }) => {
1504 final_status = RunStatus::AwaitingApproval;
1505 final_run = self
1506 .store
1507 .update_run_returning(
1508 run_id,
1509 RunUpdate {
1510 status: Some(RunStatus::AwaitingApproval),
1511 cost_usd: Some(ctx.total_cost_usd()),
1512 duration_ms: Some(total_duration),
1513 ..RunUpdate::default()
1514 },
1515 )
1516 .await?;
1517
1518 info!(
1519 run_id = %approval_run_id,
1520 step_id = %step_id,
1521 message = %message,
1522 "run awaiting approval"
1523 );
1524
1525 let requirement = self
1527 .store
1528 .get_step(step_id)
1529 .await?
1530 .and_then(|s| s.approval_requirement);
1531 self.event_publisher
1532 .publish(Event::ApprovalRequested(ApprovalRequestedEvent {
1533 run_id: approval_run_id,
1534 step_id,
1535 message: message.clone(),
1536 requirement,
1537 at: Utc::now(),
1538 }));
1539 }
1540 Err(EngineError::HumanInputRequired {
1541 run_id: input_run_id,
1542 step_id,
1543 ref message,
1544 }) => {
1545 final_status = RunStatus::AwaitingApproval;
1546 final_run = self
1547 .store
1548 .update_run_returning(
1549 run_id,
1550 RunUpdate {
1551 status: Some(RunStatus::AwaitingApproval),
1552 cost_usd: Some(ctx.total_cost_usd()),
1553 duration_ms: Some(total_duration),
1554 ..RunUpdate::default()
1555 },
1556 )
1557 .await?;
1558
1559 info!(
1561 run_id = %input_run_id,
1562 step_id = %step_id,
1563 message = %message,
1564 "run awaiting human input"
1565 );
1566 }
1567 Err(EngineError::DelaySleeping {
1568 run_id: delay_run_id,
1569 step_id,
1570 wake_at,
1571 }) => {
1572 final_status = RunStatus::Sleeping;
1573 final_run = self
1574 .store
1575 .update_run_returning(
1576 run_id,
1577 RunUpdate {
1578 status: Some(RunStatus::Sleeping),
1579 cost_usd: Some(ctx.total_cost_usd()),
1580 duration_ms: Some(total_duration),
1581 scheduled_at: Some(wake_at),
1582 ..RunUpdate::default()
1583 },
1584 )
1585 .await?;
1586
1587 info!(
1588 run_id = %delay_run_id,
1589 step_id = %step_id,
1590 wake_at = %wake_at,
1591 "run sleeping until delay elapses"
1592 );
1593 }
1594 Err(err) => {
1595 let guardrail_stop = matches!(
1599 err,
1600 EngineError::RunBudgetExceeded { .. } | EngineError::WorkflowGuardRejected(_)
1601 );
1602
1603 final_status = if guardrail_stop {
1604 if let Err(store_err) = self
1605 .store
1606 .update_run(
1607 run_id,
1608 RunUpdate {
1609 status: Some(RunStatus::Cancelled),
1610 error: Some(err.to_string()),
1611 cost_usd: Some(ctx.total_cost_usd()),
1612 duration_ms: Some(total_duration),
1613 completed_at: Some(completed_at),
1614 ..RunUpdate::default()
1615 },
1616 )
1617 .await
1618 {
1619 error!(run_id = %run_id, store_error = %store_err, "failed to persist run cancellation");
1620 }
1621 if let Err(cleanup_err) = self
1622 .fail_orphaned_steps(run_id, "run stopped: guardrail limit reached")
1623 .await
1624 {
1625 error!(run_id = %run_id, store_error = %cleanup_err, "failed to cleanup orphaned steps");
1626 }
1627 RunStatus::Cancelled
1628 } else {
1629 self.fail_or_schedule_retry(
1630 run_id,
1631 &err.to_string(),
1632 is_run_retryable(&err),
1633 Some(ctx.total_cost_usd()),
1634 Some(total_duration),
1635 )
1636 .await
1637 .unwrap_or_else(|store_err| {
1638 error!(run_id = %run_id, store_error = %store_err, "failed to persist run failure");
1639 RunStatus::Failed
1640 })
1641 };
1642
1643 if matches!(err, EngineError::RunBudgetExceeded { .. }) {
1644 self.on_run_budget_exceeded(workflow_name, run_id, &err);
1645 }
1646
1647 error!(run_id = %run_id, status = %final_status, error = %err, "run stopped");
1648
1649 self.publish_run_status_changed(
1650 workflow_name,
1651 run_id,
1652 final_status,
1653 Some(err.to_string()),
1654 ctx,
1655 total_duration,
1656 run_labels,
1657 );
1658
1659 #[cfg(feature = "prometheus")]
1660 self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
1661
1662 return Err(err);
1663 }
1664 }
1665
1666 self.publish_run_status_changed(
1667 workflow_name,
1668 run_id,
1669 final_status,
1670 None,
1671 ctx,
1672 total_duration,
1673 run_labels,
1674 );
1675
1676 #[cfg(feature = "prometheus")]
1677 self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
1678
1679 Ok(WorkflowResult {
1680 run: final_run,
1681 steps: ctx.step_results().to_vec(),
1682 })
1683 }
1684
1685 #[cfg(feature = "prometheus")]
1687 fn emit_run_metrics(
1688 &self,
1689 workflow_name: &str,
1690 status: RunStatus,
1691 duration_ms: u64,
1692 ctx: &WorkflowContext,
1693 ) {
1694 let status_str = status.to_string();
1695 let wf = workflow_name.to_string();
1696
1697 counter!(RUNS_TOTAL, "workflow" => wf.clone(), "status" => status_str.clone()).increment(1);
1698 histogram!(RUN_DURATION_SECONDS, "workflow" => wf.clone(), "status" => status_str)
1699 .record(duration_ms as f64 / 1000.0);
1700 histogram!(RUN_COST_USD, "workflow" => wf.clone()).record(
1701 ctx.total_cost_usd()
1702 .to_string()
1703 .parse::<f64>()
1704 .unwrap_or(0.0),
1705 );
1706 gauge!(RUNS_ACTIVE, "workflow" => wf).decrement(1.0);
1707 }
1708
1709 fn on_run_budget_exceeded(&self, workflow_name: &str, run_id: Uuid, err: &EngineError) {
1715 let EngineError::RunBudgetExceeded {
1716 limit_usd,
1717 spent_usd,
1718 step_budget_usd,
1719 ..
1720 } = err
1721 else {
1722 return;
1723 };
1724
1725 #[cfg(feature = "prometheus")]
1726 counter!(
1727 RUN_BUDGET_EXCEEDED_TOTAL,
1728 "workflow" => workflow_name.to_string(),
1729 "scope" => "run",
1730 )
1731 .increment(1);
1732
1733 self.event_publisher
1734 .publish(Event::RunBudgetExceeded(RunBudgetExceededEvent {
1735 run_id,
1736 workflow_name: workflow_name.to_string(),
1737 limit_usd: *limit_usd,
1738 spent_usd: *spent_usd,
1739 step_budget_usd: *step_budget_usd,
1740 at: Utc::now(),
1741 }));
1742 }
1743
1744 #[allow(clippy::too_many_arguments)]
1749 fn publish_run_status_changed(
1750 &self,
1751 workflow_name: &str,
1752 run_id: Uuid,
1753 to: RunStatus,
1754 error: Option<String>,
1755 ctx: &WorkflowContext,
1756 duration_ms: u64,
1757 labels: HashMap<String, String>,
1758 ) {
1759 let now = Utc::now();
1760 let cost_usd = ctx.total_cost_usd();
1761 let wf = workflow_name.to_string();
1762
1763 self.event_publisher
1764 .publish(Event::RunStatusChanged(RunStatusChangedEvent {
1765 run_id,
1766 workflow_name: wf.clone(),
1767 from: RunStatus::Running,
1768 to,
1769 error: error.clone(),
1770 cost_usd,
1771 duration_ms,
1772 labels: labels.clone(),
1773 at: now,
1774 }));
1775
1776 if to == RunStatus::Failed {
1777 self.event_publisher
1778 .publish(Event::RunFailed(RunFailedEvent {
1779 run_id,
1780 workflow_name: wf,
1781 error,
1782 cost_usd,
1783 duration_ms,
1784 labels,
1785 at: now,
1786 }));
1787 }
1788 }
1789}
1790
1791impl fmt::Debug for Engine {
1792 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1793 f.debug_struct("Engine")
1794 .field("handlers", &self.handlers.keys().collect::<Vec<_>>())
1795 .finish_non_exhaustive()
1796 }
1797}
1798
1799#[cfg(test)]
1800mod tests {
1801 use super::*;
1802 use crate::config::ShellConfig;
1803 use crate::handler::{HandlerFuture, WorkflowHandler};
1804 use ironflow_core::providers::claude::ClaudeCodeProvider;
1805 use ironflow_core::providers::record_replay::RecordReplayProvider;
1806 use ironflow_store::memory::InMemoryStore;
1807 use ironflow_store::models::StepStatus;
1808 use serde_json::json;
1809
1810 struct EchoWorkflow;
1812
1813 impl WorkflowHandler for EchoWorkflow {
1814 fn name(&self) -> &str {
1815 "echo-workflow"
1816 }
1817
1818 fn describe(&self) -> WorkflowInfo {
1819 WorkflowInfo {
1820 description: "A simple workflow that echoes hello".to_string(),
1821 source_code: None,
1822 sub_workflows: Vec::new(),
1823 category: None,
1824 version: self.version().map(str::to_string),
1825 compatible_versions: Vec::new(),
1826 input_schema: None,
1827 default_labels: HashMap::new(),
1828 schedule: self.schedule().cloned(),
1829 default_max_cost_usd: self.default_max_cost_usd(),
1830 }
1831 }
1832
1833 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
1834 Box::pin(async move {
1835 ctx.shell("greet", ShellConfig::new("echo hello")).await?;
1836 Ok(())
1837 })
1838 }
1839 }
1840
1841 struct FailingWorkflow;
1843
1844 impl WorkflowHandler for FailingWorkflow {
1845 fn name(&self) -> &str {
1846 "failing-workflow"
1847 }
1848
1849 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
1850 Box::pin(async move {
1851 ctx.shell("fail", ShellConfig::new("exit 1")).await?;
1852 Ok(())
1853 })
1854 }
1855 }
1856
1857 fn create_test_engine() -> Engine {
1858 let store = Arc::new(InMemoryStore::new());
1859 let inner = ClaudeCodeProvider::new();
1860 let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
1861 inner,
1862 "/tmp/ironflow-fixtures",
1863 ));
1864 Engine::new(store, provider)
1865 }
1866
1867 #[test]
1868 fn engine_new_creates_instance() {
1869 let engine = create_test_engine();
1870 assert_eq!(engine.handler_names().len(), 0);
1871 }
1872
1873 #[test]
1874 fn execution_mode_defaults_to_local() {
1875 let engine = create_test_engine();
1876 assert_eq!(engine.execution_mode(), ExecutionMode::Local);
1877 }
1878
1879 #[test]
1880 fn with_execution_mode_overrides_the_default() {
1881 let engine = create_test_engine().with_execution_mode(ExecutionMode::Workers);
1882 assert_eq!(engine.execution_mode(), ExecutionMode::Workers);
1883 }
1884
1885 #[test]
1886 fn engine_register_handler() {
1887 let mut engine = create_test_engine();
1888 let result = engine.register(EchoWorkflow);
1889 assert!(result.is_ok());
1890 assert_eq!(engine.handler_names().len(), 1);
1891 assert!(engine.handler_names().contains(&"echo-workflow"));
1892 }
1893
1894 #[test]
1895 fn engine_register_duplicate_returns_error() {
1896 let mut engine = create_test_engine();
1897 engine.register(EchoWorkflow).unwrap();
1898 let result = engine.register(EchoWorkflow);
1899 assert!(result.is_err());
1900 }
1901
1902 #[test]
1903 fn engine_get_handler_found() {
1904 let mut engine = create_test_engine();
1905 engine.register(EchoWorkflow).unwrap();
1906 let handler = engine.get_handler("echo-workflow");
1907 assert!(handler.is_some());
1908 }
1909
1910 #[test]
1911 fn engine_get_handler_not_found() {
1912 let engine = create_test_engine();
1913 let handler = engine.get_handler("nonexistent");
1914 assert!(handler.is_none());
1915 }
1916
1917 #[test]
1918 fn engine_handler_names_lists_all() {
1919 let mut engine = create_test_engine();
1920 engine.register(EchoWorkflow).unwrap();
1921 engine.register(FailingWorkflow).unwrap();
1922 let names = engine.handler_names();
1923 assert_eq!(names.len(), 2);
1924 assert!(names.contains(&"echo-workflow"));
1925 assert!(names.contains(&"failing-workflow"));
1926 }
1927
1928 #[test]
1929 fn engine_handler_info_returns_description() {
1930 let mut engine = create_test_engine();
1931 engine.register(EchoWorkflow).unwrap();
1932 let info = engine.handler_info("echo-workflow");
1933 assert!(info.is_some());
1934 let info = info.unwrap();
1935 assert_eq!(info.description, "A simple workflow that echoes hello");
1936 }
1937
1938 struct CategorizedWorkflow;
1939
1940 impl WorkflowHandler for CategorizedWorkflow {
1941 fn name(&self) -> &str {
1942 "categorized"
1943 }
1944 fn category(&self) -> Option<&str> {
1945 Some("data/etl")
1946 }
1947 fn execute<'a>(
1948 &'a self,
1949 _ctx: &'a mut WorkflowContext,
1950 ) -> crate::handler::HandlerFuture<'a> {
1951 Box::pin(async move { Ok(()) })
1952 }
1953 }
1954
1955 #[test]
1956 fn engine_default_describe_propagates_category() {
1957 let mut engine = create_test_engine();
1958 engine.register(CategorizedWorkflow).unwrap();
1959 let info = engine.handler_info("categorized").unwrap();
1960 assert_eq!(info.category.as_deref(), Some("data/etl"));
1961 }
1962
1963 #[test]
1964 fn engine_default_describe_without_category() {
1965 let mut engine = create_test_engine();
1966 engine.register(EchoWorkflow).unwrap();
1967 let info = engine.handler_info("echo-workflow").unwrap();
1968 assert!(info.category.is_none());
1969 }
1970
1971 struct ScheduledWorkflow {
1976 schedule: CronSchedule,
1977 }
1978
1979 impl ScheduledWorkflow {
1980 fn new() -> Self {
1981 Self {
1982 schedule: CronSchedule::new("0 0 * * * *").unwrap(),
1983 }
1984 }
1985 }
1986
1987 impl WorkflowHandler for ScheduledWorkflow {
1988 fn name(&self) -> &str {
1989 "scheduled"
1990 }
1991 fn schedule(&self) -> Option<&CronSchedule> {
1992 Some(&self.schedule)
1993 }
1994 fn execute<'a>(
1995 &'a self,
1996 _ctx: &'a mut WorkflowContext,
1997 ) -> crate::handler::HandlerFuture<'a> {
1998 Box::pin(async move { Ok(()) })
1999 }
2000 }
2001
2002 #[test]
2003 fn engine_default_describe_propagates_schedule() {
2004 let mut engine = create_test_engine();
2005 engine.register(ScheduledWorkflow::new()).unwrap();
2006 let info = engine.handler_info("scheduled").unwrap();
2007 assert_eq!(
2008 info.schedule.as_ref().map(|s| s.as_str()),
2009 Some("0 0 * * * *")
2010 );
2011 }
2012
2013 #[test]
2014 fn engine_default_describe_without_schedule() {
2015 let mut engine = create_test_engine();
2016 engine.register(EchoWorkflow).unwrap();
2017 let info = engine.handler_info("echo-workflow").unwrap();
2018 assert!(info.schedule.is_none());
2019 }
2020
2021 #[test]
2022 fn scheduled_handlers_returns_only_scheduled() {
2023 let mut engine = create_test_engine();
2024 engine.register(EchoWorkflow).unwrap();
2025 engine.register(ScheduledWorkflow::new()).unwrap();
2026 engine.register(FailingWorkflow).unwrap();
2027
2028 let scheduled = engine.scheduled_handlers();
2029 assert_eq!(scheduled.len(), 1);
2030 assert_eq!(scheduled[0].0, "scheduled");
2031 assert_eq!(scheduled[0].1.as_str(), "0 0 * * * *");
2032 }
2033
2034 #[test]
2035 fn scheduled_handlers_empty_when_none_scheduled() {
2036 let mut engine = create_test_engine();
2037 engine.register(EchoWorkflow).unwrap();
2038 engine.register(FailingWorkflow).unwrap();
2039
2040 let scheduled = engine.scheduled_handlers();
2041 assert!(scheduled.is_empty());
2042 }
2043
2044 struct BadCategoryWorkflow(&'static str);
2045
2046 impl WorkflowHandler for BadCategoryWorkflow {
2047 fn name(&self) -> &str {
2048 "bad-category"
2049 }
2050 fn category(&self) -> Option<&str> {
2051 Some(self.0)
2052 }
2053 fn execute<'a>(
2054 &'a self,
2055 _ctx: &'a mut WorkflowContext,
2056 ) -> crate::handler::HandlerFuture<'a> {
2057 Box::pin(async move { Ok(()) })
2058 }
2059 }
2060
2061 #[test]
2062 fn engine_register_rejects_empty_category() {
2063 let mut engine = create_test_engine();
2064 let err = engine.register(BadCategoryWorkflow("")).unwrap_err();
2065 match err {
2066 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty category")),
2067 other => panic!("expected InvalidWorkflow, got {other:?}"),
2068 }
2069 }
2070
2071 #[test]
2072 fn engine_register_rejects_leading_slash_category() {
2073 let mut engine = create_test_engine();
2074 let err = engine
2075 .register(BadCategoryWorkflow("/data/etl"))
2076 .unwrap_err();
2077 match err {
2078 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("leading '/'")),
2079 other => panic!("expected InvalidWorkflow, got {other:?}"),
2080 }
2081 }
2082
2083 #[test]
2084 fn engine_register_rejects_trailing_slash_category() {
2085 let mut engine = create_test_engine();
2086 let err = engine
2087 .register(BadCategoryWorkflow("data/etl/"))
2088 .unwrap_err();
2089 match err {
2090 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("trailing '/'")),
2091 other => panic!("expected InvalidWorkflow, got {other:?}"),
2092 }
2093 }
2094
2095 #[test]
2096 fn engine_register_rejects_double_slash_category() {
2097 let mut engine = create_test_engine();
2098 let err = engine
2099 .register(BadCategoryWorkflow("data//etl"))
2100 .unwrap_err();
2101 match err {
2102 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty segment")),
2103 other => panic!("expected InvalidWorkflow, got {other:?}"),
2104 }
2105 }
2106
2107 #[test]
2108 fn engine_register_rejects_whitespace_only_segment_category() {
2109 let mut engine = create_test_engine();
2110 let err = engine
2111 .register(BadCategoryWorkflow("data/ /etl"))
2112 .unwrap_err();
2113 match err {
2114 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("whitespace-only segment")),
2115 other => panic!("expected InvalidWorkflow, got {other:?}"),
2116 }
2117 }
2118
2119 #[test]
2120 fn engine_register_accepts_valid_nested_category() {
2121 let mut engine = create_test_engine();
2122 assert!(engine.register(CategorizedWorkflow).is_ok());
2123 }
2124
2125 #[tokio::test]
2126 async fn engine_unknown_workflow_returns_error() {
2127 let engine = create_test_engine();
2128 let result = engine
2129 .run_handler("unknown", TriggerKind::Manual, json!({}))
2130 .await;
2131 assert!(result.is_err());
2132 match result {
2133 Err(EngineError::InvalidWorkflow(msg)) => {
2134 assert!(msg.contains("no handler registered"));
2135 }
2136 _ => panic!("expected InvalidWorkflow error"),
2137 }
2138 }
2139
2140 #[tokio::test]
2141 async fn engine_enqueue_handler_creates_pending_run() {
2142 let mut engine = create_test_engine();
2143 engine.register(EchoWorkflow).unwrap();
2144
2145 let run = engine
2146 .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
2147 .await
2148 .unwrap();
2149 assert_eq!(run.status.state, RunStatus::Pending);
2150 assert_eq!(run.workflow_name, "echo-workflow");
2151 }
2152
2153 #[tokio::test]
2154 async fn enqueue_handler_leaves_the_run_unattributed() {
2155 let mut engine = create_test_engine();
2156 engine.register(EchoWorkflow).unwrap();
2157
2158 let run = engine
2159 .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
2160 .await
2161 .unwrap();
2162
2163 assert!(run.created_by.is_none());
2164 }
2165
2166 #[tokio::test]
2167 async fn enqueue_handler_with_options_records_the_author() {
2168 let mut engine = create_test_engine();
2169 engine.register(EchoWorkflow).unwrap();
2170 let actor = RunActor::User {
2171 user_id: Uuid::now_v7(),
2172 };
2173
2174 let run = engine
2175 .enqueue_handler_with_options(
2176 "echo-workflow",
2177 TriggerKind::Api,
2178 json!({}),
2179 EnqueueOptions {
2180 created_by: Some(actor.clone()),
2181 ..Default::default()
2182 },
2183 )
2184 .await
2185 .unwrap()
2186 .into_run();
2187
2188 assert_eq!(run.created_by, Some(actor));
2189 }
2190
2191 #[tokio::test]
2192 async fn enqueue_handler_with_options_accepts_no_author() {
2193 let mut engine = create_test_engine();
2194 engine.register(EchoWorkflow).unwrap();
2195
2196 let run = engine
2197 .enqueue_handler_with_options(
2198 "echo-workflow",
2199 TriggerKind::Cron {
2200 schedule: "0 * * * * *".to_string(),
2201 },
2202 json!({}),
2203 EnqueueOptions::default(),
2204 )
2205 .await
2206 .unwrap()
2207 .into_run();
2208
2209 assert!(run.created_by.is_none());
2210 }
2211
2212 #[tokio::test]
2213 async fn run_handler_leaves_the_run_unattributed() {
2214 let mut engine = create_test_engine();
2215 engine.register(EchoWorkflow).unwrap();
2216
2217 let run = engine
2218 .run_handler("echo-workflow", TriggerKind::Manual, json!({}))
2219 .await
2220 .unwrap()
2221 .run;
2222
2223 assert!(run.created_by.is_none());
2224 }
2225
2226 #[tokio::test]
2227 async fn engine_register_boxed() {
2228 let mut engine = create_test_engine();
2229 let handler: Box<dyn WorkflowHandler> = Box::new(EchoWorkflow);
2230 let result = engine.register_boxed(handler);
2231 assert!(result.is_ok());
2232 assert_eq!(engine.handler_names().len(), 1);
2233 }
2234
2235 #[tokio::test]
2236 async fn engine_store_and_provider_accessors() {
2237 let store = Arc::new(InMemoryStore::new());
2238 let inner = ClaudeCodeProvider::new();
2239 let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
2240 inner,
2241 "/tmp/ironflow-fixtures",
2242 ));
2243 let engine = Engine::new(store.clone(), provider.clone());
2244
2245 let _ = engine.store();
2247 let _ = engine.provider();
2248 }
2249
2250 use crate::operation::{Operation, OperationContext};
2255 use async_trait::async_trait;
2256 use ironflow_core::error::OperationError;
2257 use ironflow_store::models::StepKind;
2258
2259 struct FakeGitlabOp {
2260 project_id: u64,
2261 title: String,
2262 }
2263
2264 #[async_trait]
2265 impl Operation for FakeGitlabOp {
2266 fn kind(&self) -> &str {
2267 "gitlab"
2268 }
2269
2270 async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
2271 Ok(json!({
2272 "issue_id": 42,
2273 "project_id": self.project_id,
2274 "title": self.title,
2275 }))
2276 }
2277
2278 fn input(&self) -> Option<Value> {
2279 Some(json!({
2280 "project_id": self.project_id,
2281 "title": self.title,
2282 }))
2283 }
2284 }
2285
2286 struct FailingOp;
2287
2288 #[async_trait]
2289 impl Operation for FailingOp {
2290 fn kind(&self) -> &str {
2291 "broken-service"
2292 }
2293
2294 async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
2295 Err(OperationError::Http {
2296 status: None,
2297 message: "service unavailable".to_string(),
2298 })
2299 }
2300 }
2301
2302 struct OperationWorkflow;
2303
2304 impl WorkflowHandler for OperationWorkflow {
2305 fn name(&self) -> &str {
2306 "operation-workflow"
2307 }
2308
2309 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2310 Box::pin(async move {
2311 let op = FakeGitlabOp {
2312 project_id: 123,
2313 title: "Bug report".to_string(),
2314 };
2315 ctx.operation("create-issue", &op).await?;
2316 Ok(())
2317 })
2318 }
2319 }
2320
2321 struct FailingOperationWorkflow;
2322
2323 impl WorkflowHandler for FailingOperationWorkflow {
2324 fn name(&self) -> &str {
2325 "failing-operation-workflow"
2326 }
2327
2328 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2329 Box::pin(async move {
2330 ctx.operation("broken-call", &FailingOp).await?;
2331 Ok(())
2332 })
2333 }
2334 }
2335
2336 struct MixedWorkflow;
2337
2338 impl WorkflowHandler for MixedWorkflow {
2339 fn name(&self) -> &str {
2340 "mixed-workflow"
2341 }
2342
2343 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2344 Box::pin(async move {
2345 ctx.shell("build", ShellConfig::new("echo built")).await?;
2346 let op = FakeGitlabOp {
2347 project_id: 456,
2348 title: "Deploy done".to_string(),
2349 };
2350 let result = ctx.operation("notify-gitlab", &op).await?;
2351 assert_eq!(result.output["issue_id"], 42);
2352 Ok(())
2353 })
2354 }
2355 }
2356
2357 #[tokio::test]
2358 async fn operation_step_happy_path() {
2359 let mut engine = create_test_engine();
2360 engine.register(OperationWorkflow).unwrap();
2361
2362 let run = engine
2363 .run_handler("operation-workflow", TriggerKind::Manual, json!({}))
2364 .await
2365 .unwrap()
2366 .run;
2367
2368 assert_eq!(run.status.state, RunStatus::Completed);
2369
2370 let steps = engine.store().list_steps(run.id).await.unwrap();
2371
2372 assert_eq!(steps.len(), 1);
2373 assert_eq!(steps[0].name, "create-issue");
2374 assert_eq!(steps[0].kind, StepKind::Custom("gitlab".to_string()));
2375 assert_eq!(
2376 steps[0].status.state,
2377 ironflow_store::models::StepStatus::Completed
2378 );
2379
2380 let output = steps[0].output.as_ref().unwrap();
2381 assert_eq!(output["issue_id"], 42);
2382 assert_eq!(output["project_id"], 123);
2383
2384 let input = steps[0].input.as_ref().unwrap();
2385 assert_eq!(input["project_id"], 123);
2386 assert_eq!(input["title"], "Bug report");
2387 }
2388
2389 #[tokio::test]
2390 async fn operation_step_failure_marks_run_failed() {
2391 let mut engine = create_test_engine();
2392 engine.register(FailingOperationWorkflow).unwrap();
2393
2394 let result = engine
2395 .run_handler("failing-operation-workflow", TriggerKind::Manual, json!({}))
2396 .await;
2397
2398 assert!(result.is_err());
2399 }
2400
2401 #[tokio::test]
2402 async fn operation_mixed_with_shell_steps() {
2403 let mut engine = create_test_engine();
2404 engine.register(MixedWorkflow).unwrap();
2405
2406 let run = engine
2407 .run_handler("mixed-workflow", TriggerKind::Manual, json!({}))
2408 .await
2409 .unwrap()
2410 .run;
2411
2412 assert_eq!(run.status.state, RunStatus::Completed);
2413
2414 let steps = engine.store().list_steps(run.id).await.unwrap();
2415
2416 assert_eq!(steps.len(), 2);
2417 assert_eq!(steps[0].kind, StepKind::Shell);
2418 assert_eq!(steps[1].kind, StepKind::Custom("gitlab".to_string()));
2419 assert_eq!(steps[0].position, 0);
2420 assert_eq!(steps[1].position, 1);
2421 }
2422
2423 use crate::config::ApprovalConfig;
2428
2429 struct SingleApprovalWorkflow;
2430
2431 impl WorkflowHandler for SingleApprovalWorkflow {
2432 fn name(&self) -> &str {
2433 "single-approval"
2434 }
2435
2436 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2437 Box::pin(async move {
2438 ctx.shell("build", ShellConfig::new("echo built")).await?;
2439 ctx.approval("gate", ApprovalConfig::new("OK?")).await?;
2440 ctx.shell("deploy", ShellConfig::new("echo deployed"))
2441 .await?;
2442 Ok(())
2443 })
2444 }
2445 }
2446
2447 struct DoubleApprovalWorkflow;
2448
2449 impl WorkflowHandler for DoubleApprovalWorkflow {
2450 fn name(&self) -> &str {
2451 "double-approval"
2452 }
2453
2454 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2455 Box::pin(async move {
2456 ctx.shell("build", ShellConfig::new("echo built")).await?;
2457 ctx.approval("staging-gate", ApprovalConfig::new("Deploy staging?"))
2458 .await?;
2459 ctx.shell("deploy-staging", ShellConfig::new("echo staging"))
2460 .await?;
2461 ctx.approval("prod-gate", ApprovalConfig::new("Deploy prod?"))
2462 .await?;
2463 ctx.shell("deploy-prod", ShellConfig::new("echo prod"))
2464 .await?;
2465 Ok(())
2466 })
2467 }
2468 }
2469
2470 #[tokio::test]
2471 async fn approval_pauses_run() {
2472 let mut engine = create_test_engine();
2473 engine.register(SingleApprovalWorkflow).unwrap();
2474
2475 let run = engine
2476 .run_handler("single-approval", TriggerKind::Manual, json!({}))
2477 .await
2478 .unwrap()
2479 .run;
2480
2481 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
2482
2483 let steps = engine.store().list_steps(run.id).await.unwrap();
2484 assert_eq!(steps.len(), 2); assert_eq!(steps[0].kind, StepKind::Shell);
2486 assert_eq!(steps[0].status.state, StepStatus::Completed);
2487 assert_eq!(steps[1].kind, StepKind::Approval);
2488 assert_eq!(steps[1].status.state, StepStatus::AwaitingApproval);
2489 }
2490
2491 #[tokio::test]
2492 async fn approval_resume_completes_run() {
2493 let mut engine = create_test_engine();
2494 engine.register(SingleApprovalWorkflow).unwrap();
2495
2496 let run = engine
2498 .run_handler("single-approval", TriggerKind::Manual, json!({}))
2499 .await
2500 .unwrap()
2501 .run;
2502 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
2503
2504 engine
2506 .store()
2507 .update_run_status(run.id, RunStatus::Running)
2508 .await
2509 .unwrap();
2510
2511 let resumed = engine.resume_run(run.id).await.unwrap().run;
2513 assert_eq!(resumed.status.state, RunStatus::Completed);
2514
2515 let steps = engine.store().list_steps(run.id).await.unwrap();
2516 assert_eq!(steps.len(), 3); assert_eq!(steps[0].name, "build");
2518 assert_eq!(steps[0].status.state, StepStatus::Completed);
2519 assert_eq!(steps[1].name, "gate");
2520 assert_eq!(steps[1].kind, StepKind::Approval);
2521 assert_eq!(steps[1].status.state, StepStatus::Completed);
2522 assert_eq!(steps[2].name, "deploy");
2523 assert_eq!(steps[2].status.state, StepStatus::Completed);
2524 }
2525
2526 #[tokio::test]
2527 async fn double_approval_two_resumes() {
2528 let mut engine = create_test_engine();
2529 engine.register(DoubleApprovalWorkflow).unwrap();
2530
2531 let run = engine
2533 .run_handler("double-approval", TriggerKind::Manual, json!({}))
2534 .await
2535 .unwrap()
2536 .run;
2537 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
2538
2539 let steps = engine.store().list_steps(run.id).await.unwrap();
2540 assert_eq!(steps.len(), 2); engine
2544 .store()
2545 .update_run_status(run.id, RunStatus::Running)
2546 .await
2547 .unwrap();
2548
2549 let resumed = engine.resume_run(run.id).await.unwrap().run;
2550 assert_eq!(resumed.status.state, RunStatus::AwaitingApproval);
2551
2552 let steps = engine.store().list_steps(run.id).await.unwrap();
2553 assert_eq!(steps.len(), 4); engine
2557 .store()
2558 .update_run_status(run.id, RunStatus::Running)
2559 .await
2560 .unwrap();
2561
2562 let final_run = engine.resume_run(run.id).await.unwrap().run;
2563 assert_eq!(final_run.status.state, RunStatus::Completed);
2564
2565 let steps = engine.store().list_steps(run.id).await.unwrap();
2566 assert_eq!(steps.len(), 5);
2567 assert_eq!(steps[0].name, "build");
2568 assert_eq!(steps[1].name, "staging-gate");
2569 assert_eq!(steps[2].name, "deploy-staging");
2570 assert_eq!(steps[3].name, "prod-gate");
2571 assert_eq!(steps[4].name, "deploy-prod");
2572
2573 for step in &steps {
2574 assert_eq!(step.status.state, StepStatus::Completed);
2575 }
2576 }
2577
2578 use ironflow_store::models::{NewStep, StepUpdate, step_trace_id};
2583
2584 async fn create_step_with_status(
2585 store: &Arc<dyn Store>,
2586 run_id: Uuid,
2587 name: &str,
2588 position: u32,
2589 status: StepStatus,
2590 ) -> ironflow_store::models::Step {
2591 let step = store
2592 .create_step(NewStep {
2593 run_id,
2594 trace_id: step_trace_id(run_id, name, position),
2595 name: name.to_string(),
2596 kind: StepKind::Shell,
2597 position,
2598 input: None,
2599 is_error_handler: false,
2600 })
2601 .await
2602 .unwrap();
2603
2604 match status {
2605 StepStatus::Pending => {}
2606 StepStatus::Running => {
2607 store
2608 .update_step(
2609 step.id,
2610 StepUpdate {
2611 status: Some(StepStatus::Running),
2612 ..StepUpdate::default()
2613 },
2614 )
2615 .await
2616 .unwrap();
2617 }
2618 StepStatus::Completed => {
2619 store
2620 .update_step(
2621 step.id,
2622 StepUpdate {
2623 status: Some(StepStatus::Running),
2624 ..StepUpdate::default()
2625 },
2626 )
2627 .await
2628 .unwrap();
2629 store
2630 .update_step(
2631 step.id,
2632 StepUpdate {
2633 status: Some(StepStatus::Completed),
2634 ..StepUpdate::default()
2635 },
2636 )
2637 .await
2638 .unwrap();
2639 }
2640 StepStatus::AwaitingApproval => {
2641 store
2642 .update_step(
2643 step.id,
2644 StepUpdate {
2645 status: Some(StepStatus::Running),
2646 ..StepUpdate::default()
2647 },
2648 )
2649 .await
2650 .unwrap();
2651 store
2652 .update_step(
2653 step.id,
2654 StepUpdate {
2655 status: Some(StepStatus::AwaitingApproval),
2656 ..StepUpdate::default()
2657 },
2658 )
2659 .await
2660 .unwrap();
2661 }
2662 _ => panic!("unsupported status for test helper: {status}"),
2663 }
2664
2665 store.get_step(step.id).await.unwrap().unwrap()
2666 }
2667
2668 #[tokio::test]
2669 async fn fail_orphaned_steps_marks_running_as_failed() {
2670 let engine = create_test_engine();
2671 let run = engine
2672 .store()
2673 .create_run(NewRun {
2674 created_by: None,
2675 workflow_name: "test".to_string(),
2676 trigger: TriggerKind::Manual,
2677 payload: json!({}),
2678 max_retries: 0,
2679 handler_version: None,
2680 labels: HashMap::new(),
2681 scheduled_at: None,
2682 idempotency_key: None,
2683 max_cost_usd: None,
2684 })
2685 .await
2686 .unwrap()
2687 .into_run();
2688
2689 let step = create_step_with_status(
2690 engine.store(),
2691 run.id,
2692 "running-step",
2693 0,
2694 StepStatus::Running,
2695 )
2696 .await;
2697
2698 engine
2699 .fail_orphaned_steps(run.id, "parent run timed out")
2700 .await
2701 .unwrap();
2702
2703 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
2704 assert_eq!(updated.status.state, StepStatus::Failed);
2705 assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
2706 assert!(updated.completed_at.is_some());
2707 }
2708
2709 #[tokio::test]
2710 async fn fail_orphaned_steps_marks_pending_as_skipped() {
2711 let engine = create_test_engine();
2712 let run = engine
2713 .store()
2714 .create_run(NewRun {
2715 created_by: None,
2716 workflow_name: "test".to_string(),
2717 trigger: TriggerKind::Manual,
2718 payload: json!({}),
2719 max_retries: 0,
2720 handler_version: None,
2721 labels: HashMap::new(),
2722 scheduled_at: None,
2723 idempotency_key: None,
2724 max_cost_usd: None,
2725 })
2726 .await
2727 .unwrap()
2728 .into_run();
2729
2730 let step = create_step_with_status(
2731 engine.store(),
2732 run.id,
2733 "pending-step",
2734 0,
2735 StepStatus::Pending,
2736 )
2737 .await;
2738
2739 engine
2740 .fail_orphaned_steps(run.id, "parent run timed out")
2741 .await
2742 .unwrap();
2743
2744 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
2745 assert_eq!(updated.status.state, StepStatus::Skipped);
2746 assert!(updated.error.is_none());
2747 assert!(updated.completed_at.is_some());
2748 }
2749
2750 #[tokio::test]
2751 async fn fail_orphaned_steps_marks_awaiting_approval_as_failed() {
2752 let engine = create_test_engine();
2753 let run = engine
2754 .store()
2755 .create_run(NewRun {
2756 created_by: None,
2757 workflow_name: "test".to_string(),
2758 trigger: TriggerKind::Manual,
2759 payload: json!({}),
2760 max_retries: 0,
2761 handler_version: None,
2762 labels: HashMap::new(),
2763 scheduled_at: None,
2764 idempotency_key: None,
2765 max_cost_usd: None,
2766 })
2767 .await
2768 .unwrap()
2769 .into_run();
2770
2771 let step = create_step_with_status(
2772 engine.store(),
2773 run.id,
2774 "approval-step",
2775 0,
2776 StepStatus::AwaitingApproval,
2777 )
2778 .await;
2779
2780 engine
2781 .fail_orphaned_steps(run.id, "parent run timed out")
2782 .await
2783 .unwrap();
2784
2785 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
2786 assert_eq!(updated.status.state, StepStatus::Failed);
2787 assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
2788 assert!(updated.completed_at.is_some());
2789 }
2790
2791 #[tokio::test]
2792 async fn fail_orphaned_steps_skips_terminal_steps() {
2793 let engine = create_test_engine();
2794 let run = engine
2795 .store()
2796 .create_run(NewRun {
2797 created_by: None,
2798 workflow_name: "test".to_string(),
2799 trigger: TriggerKind::Manual,
2800 payload: json!({}),
2801 max_retries: 0,
2802 handler_version: None,
2803 labels: HashMap::new(),
2804 scheduled_at: None,
2805 idempotency_key: None,
2806 max_cost_usd: None,
2807 })
2808 .await
2809 .unwrap()
2810 .into_run();
2811
2812 let completed_step =
2813 create_step_with_status(engine.store(), run.id, "done", 0, StepStatus::Completed).await;
2814 let running_step =
2815 create_step_with_status(engine.store(), run.id, "in-flight", 1, StepStatus::Running)
2816 .await;
2817
2818 engine
2819 .fail_orphaned_steps(run.id, "parent run timed out")
2820 .await
2821 .unwrap();
2822
2823 let completed = engine
2824 .store()
2825 .get_step(completed_step.id)
2826 .await
2827 .unwrap()
2828 .unwrap();
2829 assert_eq!(completed.status.state, StepStatus::Completed);
2830
2831 let failed = engine
2832 .store()
2833 .get_step(running_step.id)
2834 .await
2835 .unwrap()
2836 .unwrap();
2837 assert_eq!(failed.status.state, StepStatus::Failed);
2838 }
2839
2840 #[tokio::test]
2841 async fn fail_orphaned_steps_mixed_states() {
2842 let engine = create_test_engine();
2843 let run = engine
2844 .store()
2845 .create_run(NewRun {
2846 created_by: None,
2847 workflow_name: "test".to_string(),
2848 trigger: TriggerKind::Manual,
2849 payload: json!({}),
2850 max_retries: 0,
2851 handler_version: None,
2852 labels: HashMap::new(),
2853 scheduled_at: None,
2854 idempotency_key: None,
2855 max_cost_usd: None,
2856 })
2857 .await
2858 .unwrap()
2859 .into_run();
2860
2861 let s_completed =
2862 create_step_with_status(engine.store(), run.id, "step-1", 0, StepStatus::Completed)
2863 .await;
2864 let s_running =
2865 create_step_with_status(engine.store(), run.id, "step-2", 1, StepStatus::Running).await;
2866 let s_pending =
2867 create_step_with_status(engine.store(), run.id, "step-3", 2, StepStatus::Pending).await;
2868
2869 engine.fail_orphaned_steps(run.id, "timeout").await.unwrap();
2870
2871 let r_completed = engine
2872 .store()
2873 .get_step(s_completed.id)
2874 .await
2875 .unwrap()
2876 .unwrap();
2877 assert_eq!(r_completed.status.state, StepStatus::Completed);
2878
2879 let r_running = engine
2880 .store()
2881 .get_step(s_running.id)
2882 .await
2883 .unwrap()
2884 .unwrap();
2885 assert_eq!(r_running.status.state, StepStatus::Failed);
2886 assert_eq!(r_running.error.as_deref(), Some("timeout"));
2887
2888 let r_pending = engine
2889 .store()
2890 .get_step(s_pending.id)
2891 .await
2892 .unwrap()
2893 .unwrap();
2894 assert_eq!(r_pending.status.state, StepStatus::Skipped);
2895 assert!(r_pending.error.is_none());
2896 }
2897
2898 #[tokio::test]
2899 async fn fail_orphaned_steps_no_steps_is_noop() {
2900 let engine = create_test_engine();
2901 let run = engine
2902 .store()
2903 .create_run(NewRun {
2904 created_by: None,
2905 workflow_name: "test".to_string(),
2906 trigger: TriggerKind::Manual,
2907 payload: json!({}),
2908 max_retries: 0,
2909 handler_version: None,
2910 labels: HashMap::new(),
2911 scheduled_at: None,
2912 idempotency_key: None,
2913 max_cost_usd: None,
2914 })
2915 .await
2916 .unwrap()
2917 .into_run();
2918
2919 let result = engine.fail_orphaned_steps(run.id, "timeout").await;
2920 assert!(result.is_ok());
2921 }
2922
2923 #[tokio::test]
2924 async fn fail_orphaned_steps_preserves_existing_error() {
2925 let engine = create_test_engine();
2926 let run = engine
2927 .store()
2928 .create_run(NewRun {
2929 created_by: None,
2930 workflow_name: "test".to_string(),
2931 trigger: TriggerKind::Manual,
2932 payload: json!({}),
2933 max_retries: 0,
2934 handler_version: None,
2935 labels: HashMap::new(),
2936 scheduled_at: None,
2937 idempotency_key: None,
2938 max_cost_usd: None,
2939 })
2940 .await
2941 .unwrap()
2942 .into_run();
2943
2944 let step_with_error = create_step_with_status(
2945 engine.store(),
2946 run.id,
2947 "already-errored",
2948 0,
2949 StepStatus::Running,
2950 )
2951 .await;
2952
2953 engine
2954 .store()
2955 .update_step(
2956 step_with_error.id,
2957 StepUpdate {
2958 error: Some("real error from provider".to_string()),
2959 ..StepUpdate::default()
2960 },
2961 )
2962 .await
2963 .unwrap();
2964
2965 let step_no_error = create_step_with_status(
2966 engine.store(),
2967 run.id,
2968 "no-error-yet",
2969 1,
2970 StepStatus::Running,
2971 )
2972 .await;
2973
2974 engine
2975 .fail_orphaned_steps(run.id, "parent run failed")
2976 .await
2977 .unwrap();
2978
2979 let updated_with = engine
2980 .store()
2981 .get_step(step_with_error.id)
2982 .await
2983 .unwrap()
2984 .unwrap();
2985 assert_eq!(updated_with.status.state, StepStatus::Failed);
2986 assert_eq!(
2987 updated_with.error.as_deref(),
2988 Some("real error from provider"),
2989 );
2990
2991 let updated_without = engine
2992 .store()
2993 .get_step(step_no_error.id)
2994 .await
2995 .unwrap()
2996 .unwrap();
2997 assert_eq!(updated_without.status.state, StepStatus::Failed);
2998 assert_eq!(updated_without.error.as_deref(), Some("parent run failed"),);
2999 }
3000}