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, 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;
27use ironflow_store::error::StoreError;
28use ironflow_store::models::{
29 NewRun, NewSignal, Run, RunActor, RunCreation, RunFilter, RunStatus, RunUpdate, SignalInsert,
30 SignalStepResolution, StepStatus, StepUpdate, TriggerKind,
31};
32use ironflow_store::store::Store;
33#[cfg(feature = "prometheus")]
34use metrics::{counter, gauge, histogram};
35
36use crate::artifact::ArtifactSink;
37use crate::budget::{BudgetConfig, month_start};
38use crate::context::WorkflowContext;
39use crate::error::EngineError;
40use crate::executor::{StepInterceptor, StepResult};
41use crate::guard::{WorkflowGuardConfig, new_shared_guard_state};
42use crate::handler::{WorkflowHandler, WorkflowInfo};
43use crate::log_sender::LogSender;
44use crate::notify::{
45 ApprovalRequestedEvent, Event, EventPublisher, EventSubscriber, RunBudgetExceededEvent,
46 RunFailedEvent, RunStatusChangedEvent, SignalAwaitedEvent, SignalReceivedEvent,
47 WorkflowEventBus,
48};
49use crate::plan::{
50 ExecutionPlan, PlanOptions, PlanRecorder, SharedPlanRecorder, estimate_durations, lock_plan,
51};
52use crate::retry_policy::{backoff_for_retry, is_run_retryable};
53use crate::schedule::CronSchedule;
54use crate::signal::{
55 Signal, SignalDelivery, SignalRejected, SignalResumed, received_output, validate_step_payload,
56};
57use ironflow_core::decision::DecisionProvider;
58
59#[derive(Debug, Clone)]
78pub struct WorkflowResult {
79 pub run: Run,
81 pub steps: Vec<StepResult>,
83}
84
85#[derive(Debug, Clone, Default)]
104pub struct EnqueueOptions {
105 pub max_retries: u32,
107 pub labels: HashMap<String, String>,
109 pub scheduled_at: Option<DateTime<Utc>>,
112 pub max_cost_usd: Option<Decimal>,
116 pub created_by: Option<RunActor>,
119 pub idempotency_key: Option<String>,
125}
126
127#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
138pub enum ExecutionMode {
139 #[default]
143 Local,
144 Workers,
147}
148
149pub struct Engine {
190 store: Arc<dyn Store>,
191 provider: Arc<dyn AgentProvider>,
192 handlers: HashMap<String, Arc<dyn WorkflowHandler>>,
193 event_publisher: EventPublisher,
194 log_sender: Option<LogSender>,
195 budget: BudgetConfig,
196 artifact_sink: Option<Arc<dyn ArtifactSink>>,
197 guard_config: Option<WorkflowGuardConfig>,
198 event_bus: Option<WorkflowEventBus>,
199 decision_provider: Option<Arc<dyn DecisionProvider>>,
200 step_interceptor: Option<Arc<dyn StepInterceptor>>,
201 execution_mode: ExecutionMode,
202}
203
204fn validate_category(handler_name: &str, category: &str) -> Result<(), EngineError> {
214 let reject = |reason: &str| {
215 Err(EngineError::InvalidWorkflow(format!(
216 "handler '{handler_name}' has invalid category '{category}': {reason}"
217 )))
218 };
219
220 if category.is_empty() {
221 return reject("empty category");
222 }
223 if category.starts_with('/') {
224 return reject("leading '/'");
225 }
226 if category.ends_with('/') {
227 return reject("trailing '/'");
228 }
229 for segment in category.split('/') {
230 if segment.is_empty() {
231 return reject("empty segment (double '/')");
232 }
233 if segment.trim().is_empty() {
234 return reject("whitespace-only segment");
235 }
236 }
237 Ok(())
238}
239
240impl Engine {
241 pub fn new(store: Arc<dyn Store>, provider: Arc<dyn AgentProvider>) -> Self {
257 Self {
258 store,
259 provider,
260 handlers: HashMap::new(),
261 event_publisher: EventPublisher::new(),
262 log_sender: None,
263 budget: BudgetConfig::new(),
264 artifact_sink: None,
265 guard_config: None,
266 event_bus: None,
267 decision_provider: None,
268 step_interceptor: None,
269 execution_mode: ExecutionMode::default(),
270 }
271 }
272
273 pub fn with_decision_provider(mut self, provider: Arc<dyn DecisionProvider>) -> Self {
294 self.decision_provider = Some(provider);
295 self
296 }
297
298 pub fn with_step_interceptor(mut self, interceptor: Arc<dyn StepInterceptor>) -> Self {
322 self.step_interceptor = Some(interceptor);
323 self
324 }
325
326 pub fn step_interceptor(&self) -> Option<&Arc<dyn StepInterceptor>> {
328 self.step_interceptor.as_ref()
329 }
330
331 pub fn with_budget_config(mut self, budget: BudgetConfig) -> Self {
352 self.budget = budget;
353 self
354 }
355
356 pub fn budget_config(&self) -> &BudgetConfig {
358 &self.budget
359 }
360
361 pub fn with_guard_config(mut self, config: WorkflowGuardConfig) -> Self {
383 self.guard_config = Some(config);
384 self
385 }
386
387 pub fn guard_config(&self) -> Option<&WorkflowGuardConfig> {
389 self.guard_config.as_ref()
390 }
391
392 pub fn with_execution_mode(mut self, mode: ExecutionMode) -> Self {
413 self.execution_mode = mode;
414 self
415 }
416
417 pub fn execution_mode(&self) -> ExecutionMode {
419 self.execution_mode
420 }
421
422 pub fn set_log_sender(&mut self, sender: LogSender) {
428 self.log_sender = Some(sender);
429 }
430
431 pub fn set_artifact_sink(&mut self, sink: Arc<dyn ArtifactSink>) {
449 self.artifact_sink = Some(sink);
450 }
451
452 pub fn artifact_sink(&self) -> Option<&Arc<dyn ArtifactSink>> {
454 self.artifact_sink.as_ref()
455 }
456
457 pub fn set_event_bus(&mut self, bus: WorkflowEventBus) {
474 self.event_bus = Some(bus);
475 }
476
477 pub fn event_bus(&self) -> Option<&WorkflowEventBus> {
479 self.event_bus.as_ref()
480 }
481
482 pub fn store(&self) -> &Arc<dyn Store> {
484 &self.store
485 }
486
487 pub fn provider(&self) -> &Arc<dyn AgentProvider> {
489 &self.provider
490 }
491
492 fn build_context(&self, run: &Run) -> WorkflowContext {
501 let handlers = self.handlers.clone();
502 let resolver: crate::context::HandlerResolver =
503 Arc::new(move |name: &str| handlers.get(name).cloned());
504 let mut ctx = WorkflowContext::with_handler_resolver(
505 run.id,
506 run.workflow_name.clone(),
507 self.store.clone(),
508 self.provider.clone(),
509 resolver,
510 );
511 ctx.carry_over_run_totals(run.retry_count + 1, run.cost_usd, run.duration_ms);
512 ctx.set_max_cost_usd(run.max_cost_usd);
513 ctx.set_run_created_at(run.created_at);
514 if let Some(ref sender) = self.log_sender {
515 ctx.set_log_sender(sender.clone());
516 }
517 if let Some(ref sink) = self.artifact_sink {
518 ctx.set_artifact_sink(sink.clone());
519 }
520 if let Some(ref bus) = self.event_bus {
521 ctx.set_event_bus(bus.clone());
522 }
523 if let Some(ref provider) = self.decision_provider {
524 ctx.set_decision_provider(provider.clone());
525 }
526 if let Some(ref interceptor) = self.step_interceptor {
527 ctx.set_step_interceptor(interceptor.clone());
528 }
529 ctx
530 }
531
532 fn build_context_with_guard(
538 &self,
539 run: &Run,
540 handler: &dyn WorkflowHandler,
541 ) -> WorkflowContext {
542 let mut ctx = self.build_context(run);
543 let guard_config = handler.guard_config().or_else(|| self.guard_config.clone());
544 if let Some(config) = guard_config {
545 ctx.set_guard(config, new_shared_guard_state());
546 }
547 ctx
548 }
549
550 async fn check_monthly_quota(&self, workflow_name: &str) -> Result<(), EngineError> {
561 let Some(limit) = self.budget.monthly_cost_limit_usd else {
562 return Ok(());
563 };
564
565 let stats = self
566 .store
567 .get_stats(RunFilter {
568 created_after: Some(month_start(Utc::now())),
569 ..RunFilter::default()
570 })
571 .await?;
572
573 if stats.total_cost_usd < limit {
574 return Ok(());
575 }
576
577 warn!(
578 workflow = %workflow_name,
579 limit_usd = %limit,
580 spent_usd = %stats.total_cost_usd,
581 "monthly cost quota exhausted, refusing new run"
582 );
583
584 #[cfg(feature = "prometheus")]
585 counter!(
586 RUN_BUDGET_EXCEEDED_TOTAL,
587 "workflow" => workflow_name.to_string(),
588 "scope" => "monthly",
589 )
590 .increment(1);
591
592 Err(EngineError::MonthlyBudgetExceeded {
593 limit_usd: limit,
594 spent_usd: stats.total_cost_usd,
595 })
596 }
597
598 pub fn register(&mut self, handler: impl WorkflowHandler + 'static) -> Result<(), EngineError> {
642 let name = handler.name().to_string();
643 if self.handlers.contains_key(&name) {
644 return Err(EngineError::InvalidWorkflow(format!(
645 "handler '{}' already registered",
646 name
647 )));
648 }
649 if let Some(category) = handler.category() {
650 validate_category(&name, category)?;
651 }
652 self.handlers.insert(name, Arc::new(handler));
653 Ok(())
654 }
655
656 pub fn register_boxed(&mut self, handler: Box<dyn WorkflowHandler>) -> Result<(), EngineError> {
663 let name = handler.name().to_string();
664 if self.handlers.contains_key(&name) {
665 return Err(EngineError::InvalidWorkflow(format!(
666 "handler '{}' already registered",
667 name
668 )));
669 }
670 if let Some(category) = handler.category() {
671 validate_category(&name, category)?;
672 }
673 self.handlers.insert(name, Arc::from(handler));
674 Ok(())
675 }
676
677 pub fn get_handler(&self, name: &str) -> Option<&Arc<dyn WorkflowHandler>> {
679 self.handlers.get(name)
680 }
681
682 pub fn handler_names(&self) -> Vec<&str> {
684 self.handlers.keys().map(|s| s.as_str()).collect()
685 }
686
687 pub fn handler_info(&self, name: &str) -> Option<WorkflowInfo> {
689 self.handlers.get(name).map(|h| h.describe())
690 }
691
692 pub fn scheduled_handlers(&self) -> Vec<(&str, &CronSchedule)> {
716 self.handlers
717 .iter()
718 .filter_map(|(name, handler)| handler.schedule().map(|sched| (name.as_str(), sched)))
719 .collect()
720 }
721
722 pub fn subscribe(
747 &mut self,
748 subscriber: impl EventSubscriber + 'static,
749 event_types: &[&'static str],
750 ) {
751 self.event_publisher.subscribe(subscriber, event_types);
752 }
753
754 pub fn event_publisher(&self) -> &EventPublisher {
759 &self.event_publisher
760 }
761
762 #[tracing::instrument(name = "engine.run_handler", skip_all, fields(workflow = %handler_name))]
792 pub async fn run_handler(
793 &self,
794 handler_name: &str,
795 trigger: TriggerKind,
796 payload: Value,
797 ) -> Result<WorkflowResult, EngineError> {
798 let handler = self
799 .handlers
800 .get(handler_name)
801 .ok_or_else(|| {
802 EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
803 })?
804 .clone();
805
806 self.check_monthly_quota(handler_name).await?;
807
808 let handler_version = handler.version().map(str::to_string);
809 let max_cost_usd = self
810 .budget
811 .resolve_run_cap(None, handler.default_max_cost_usd());
812 let run = self
813 .store
814 .create_run(NewRun {
815 created_by: None,
816 workflow_name: handler_name.to_string(),
817 trigger,
818 payload,
819 max_retries: 0,
820 handler_version,
821 labels: handler.default_labels(),
822 scheduled_at: None,
823 idempotency_key: None,
824 max_cost_usd,
825 })
826 .await?
827 .into_run();
828
829 let run_id = run.id;
830 info!(run_id = %run_id, handler_version = run.handler_version.as_deref().unwrap_or(""), "run created");
831
832 self.store
833 .update_run_status(run_id, RunStatus::Running)
834 .await?;
835
836 #[cfg(feature = "prometheus")]
837 gauge!(RUNS_ACTIVE, "workflow" => handler_name.to_string()).increment(1.0);
838
839 let run_start = Instant::now();
840 let mut ctx = self.build_context_with_guard(&run, handler.as_ref());
841
842 let result = handler.execute(&mut ctx).await;
843 self.finalize_run(run_id, handler_name, result, &ctx, run_start, run.labels)
844 .await
845 }
846
847 #[tracing::instrument(name = "engine.plan_handler", skip_all, fields(workflow = %handler_name))]
889 pub async fn plan_handler(
890 &self,
891 handler_name: &str,
892 payload: Value,
893 options: PlanOptions,
894 ) -> Result<ExecutionPlan, EngineError> {
895 if options.max_depth == 0 {
896 return Err(EngineError::InvalidWorkflow(
897 "max_depth must be at least 1".to_string(),
898 ));
899 }
900
901 let handler = self
902 .handlers
903 .get(handler_name)
904 .ok_or_else(|| {
905 EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
906 })?
907 .clone();
908
909 let estimates = if options.estimate_durations {
910 estimate_durations(&self.store, handler_name, options.sample_runs).await?
911 } else {
912 HashMap::new()
913 };
914
915 let shared: SharedPlanRecorder = Arc::new(Mutex::new(PlanRecorder::new(
916 handler_name.to_string(),
917 payload,
918 options.max_depth,
919 estimates,
920 )));
921
922 let handlers = self.handlers.clone();
925 let resolver: crate::context::HandlerResolver =
926 Arc::new(move |name: &str| handlers.get(name).cloned());
927 let mut ctx = WorkflowContext::with_handler_resolver(
928 Uuid::now_v7(),
929 handler_name.to_string(),
930 self.store.clone(),
931 self.provider.clone(),
932 resolver,
933 );
934 ctx.set_plan(shared.clone());
935
936 if let Err(err) = handler.execute(&mut ctx).await {
937 lock_plan(&shared).fail(err.to_string());
938 }
939 drop(ctx);
940
941 let plan = match Arc::try_unwrap(shared) {
942 Ok(mutex) => mutex
943 .into_inner()
944 .unwrap_or_else(|poisoned| poisoned.into_inner())
945 .into_plan(),
946 Err(shared) => lock_plan(&shared).snapshot(),
947 };
948
949 info!(
950 workflow = %handler_name,
951 steps = plan.steps.len(),
952 truncated = plan.truncated,
953 "execution plan built"
954 );
955
956 Ok(plan)
957 }
958
959 #[tracing::instrument(name = "engine.enqueue_handler", skip_all, fields(workflow = %handler_name))]
970 pub async fn enqueue_handler(
971 &self,
972 handler_name: &str,
973 trigger: TriggerKind,
974 payload: Value,
975 max_retries: u32,
976 ) -> Result<Run, EngineError> {
977 self.enqueue_handler_with_options(
978 handler_name,
979 trigger,
980 payload,
981 EnqueueOptions {
982 max_retries,
983 ..Default::default()
984 },
985 )
986 .await
987 .map(RunCreation::into_run)
988 }
989
990 #[tracing::instrument(name = "engine.enqueue_handler_with_options", skip_all, fields(workflow = %handler_name))]
1035 pub async fn enqueue_handler_with_options(
1036 &self,
1037 handler_name: &str,
1038 trigger: TriggerKind,
1039 payload: Value,
1040 options: EnqueueOptions,
1041 ) -> Result<RunCreation, EngineError> {
1042 let EnqueueOptions {
1043 max_retries,
1044 labels,
1045 scheduled_at,
1046 max_cost_usd,
1047 created_by,
1048 idempotency_key,
1049 } = options;
1050
1051 let handler = self.handlers.get(handler_name).ok_or_else(|| {
1052 EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
1053 })?;
1054
1055 self.check_monthly_quota(handler_name).await?;
1056
1057 let handler_version = handler.version().map(str::to_string);
1058 let mut merged_labels = handler.default_labels();
1059 merged_labels.extend(labels);
1060 let resolved_cap = self
1061 .budget
1062 .resolve_run_cap(max_cost_usd, handler.default_max_cost_usd());
1063
1064 let creation = self
1065 .store
1066 .create_run(NewRun {
1067 workflow_name: handler_name.to_string(),
1068 trigger,
1069 payload,
1070 max_retries,
1071 handler_version,
1072 labels: merged_labels,
1073 scheduled_at,
1074 created_by,
1075 idempotency_key,
1076 max_cost_usd: resolved_cap,
1077 })
1078 .await?;
1079
1080 match &creation {
1081 RunCreation::Created(run) => info!(
1082 run_id = %run.id,
1083 workflow = %handler_name,
1084 max_cost_usd = ?resolved_cap,
1085 "handler run enqueued"
1086 ),
1087 RunCreation::Existing(run) => info!(
1088 run_id = %run.id,
1089 workflow = %handler_name,
1090 "idempotent replay, nothing enqueued"
1091 ),
1092 }
1093
1094 Ok(creation)
1095 }
1096
1097 #[tracing::instrument(name = "engine.execute_handler_run", skip_all, fields(run_id = %run_id))]
1113 pub async fn execute_handler_run(&self, run_id: Uuid) -> Result<WorkflowResult, EngineError> {
1114 let run = self
1115 .store
1116 .get_run(run_id)
1117 .await?
1118 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
1119
1120 let handler = self
1121 .handlers
1122 .get(&run.workflow_name)
1123 .ok_or_else(|| {
1124 EngineError::InvalidWorkflow(format!(
1125 "no handler registered: {}",
1126 run.workflow_name
1127 ))
1128 })?
1129 .clone();
1130
1131 #[cfg(feature = "prometheus")]
1132 gauge!(RUNS_ACTIVE, "workflow" => run.workflow_name.clone()).increment(1.0);
1133
1134 let run_start = Instant::now();
1135 let mut ctx = self.build_context_with_guard(&run, handler.as_ref());
1136
1137 let result = if handler.is_version_compatible(run.handler_version.as_deref()) {
1148 ctx.load_replay_steps().await?;
1149 self.release_then_execute(run_id, handler.as_ref(), &mut ctx)
1150 .await
1151 } else {
1152 Err(EngineError::HandlerVersionMismatch {
1153 run_id,
1154 workflow_name: run.workflow_name.clone(),
1155 run_version: run
1156 .handler_version
1157 .clone()
1158 .unwrap_or_else(|| "unknown".to_string()),
1159 current_version: handler
1160 .version()
1161 .map(str::to_string)
1162 .unwrap_or_else(|| "unknown".to_string()),
1163 })
1164 };
1165
1166 self.finalize_run(
1167 run_id,
1168 &run.workflow_name,
1169 result,
1170 &ctx,
1171 run_start,
1172 run.labels,
1173 )
1174 .await
1175 }
1176
1177 #[tracing::instrument(name = "engine.execute_run", skip_all, fields(run_id = %run_id))]
1185 pub async fn execute_run(&self, run_id: Uuid) -> Result<WorkflowResult, EngineError> {
1186 self.execute_handler_run(run_id).await
1187 }
1188
1189 #[tracing::instrument(name = "engine.resume_run", skip_all, fields(run_id = %run_id))]
1211 pub async fn resume_run(&self, run_id: Uuid) -> Result<WorkflowResult, EngineError> {
1212 let run = self
1213 .store
1214 .get_run(run_id)
1215 .await?
1216 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
1217
1218 let handler = self
1219 .handlers
1220 .get(&run.workflow_name)
1221 .ok_or_else(|| {
1222 EngineError::InvalidWorkflow(format!(
1223 "no handler registered: {}",
1224 run.workflow_name
1225 ))
1226 })?
1227 .clone();
1228
1229 info!(run_id = %run_id, workflow = %run.workflow_name, "resuming run after approval");
1230
1231 let run_start = Instant::now();
1232 let mut ctx = self.build_context_with_guard(&run, handler.as_ref());
1233
1234 let result = if handler.is_version_compatible(run.handler_version.as_deref()) {
1235 ctx.load_replay_steps().await?;
1236 self.release_then_execute(run_id, handler.as_ref(), &mut ctx)
1237 .await
1238 } else {
1239 Err(EngineError::HandlerVersionMismatch {
1240 run_id,
1241 workflow_name: run.workflow_name.clone(),
1242 run_version: run
1243 .handler_version
1244 .clone()
1245 .unwrap_or_else(|| "unknown".to_string()),
1246 current_version: handler
1247 .version()
1248 .map(str::to_string)
1249 .unwrap_or_else(|| "unknown".to_string()),
1250 })
1251 };
1252
1253 self.finalize_run(
1254 run_id,
1255 &run.workflow_name,
1256 result,
1257 &ctx,
1258 run_start,
1259 run.labels,
1260 )
1261 .await
1262 }
1263
1264 pub async fn deliver_signal(
1310 self: &Arc<Self>,
1311 signal: NewSignal,
1312 ) -> Result<SignalDelivery, EngineError> {
1313 if signal.name.trim().is_empty() {
1314 return Err(EngineError::InvalidSignal(
1315 "signal name must not be empty".to_string(),
1316 ));
1317 }
1318 if signal.key.trim().is_empty() {
1319 return Err(EngineError::InvalidSignal(
1320 "signal key must not be empty".to_string(),
1321 ));
1322 }
1323
1324 let stored = match self.store.insert_signal(signal).await? {
1325 SignalInsert::Created(stored) => stored,
1326 SignalInsert::Duplicate(existing) => {
1327 info!(
1328 signal_id = %existing.id,
1329 signal = %existing.name,
1330 key = %existing.key,
1331 "duplicate signal ignored"
1332 );
1333 return Ok(SignalDelivery {
1334 signal_id: existing.id,
1335 duplicate: true,
1336 resumed: Vec::new(),
1337 rejected: Vec::new(),
1338 });
1339 }
1340 };
1341
1342 let waiters = self
1343 .store
1344 .list_signal_waiters(&stored.name, &stored.key)
1345 .await?;
1346 let mut resumed = Vec::new();
1347 let mut rejected = Vec::new();
1348
1349 for step in waiters {
1350 if let Err(error) = validate_step_payload(step.input.as_ref(), &stored.payload) {
1351 rejected.push(SignalRejected {
1352 run_id: step.run_id,
1353 step_id: step.id,
1354 error,
1355 });
1356 continue;
1357 }
1358
1359 match self
1360 .store
1361 .resolve_signal_step(step.id, received_output(&stored))
1362 .await
1363 {
1364 Ok(SignalStepResolution::Resolved {
1365 run_id,
1366 run_resumed,
1367 }) => {
1368 resumed.push(SignalResumed {
1369 run_id,
1370 step_id: step.id,
1371 });
1372 if run_resumed && self.execution_mode == ExecutionMode::Local {
1373 self.spawn_local_resume(run_id);
1374 }
1375 }
1376 Ok(SignalStepResolution::NotWaiting { .. }) => {}
1378 Err(err) => {
1379 error!(
1380 run_id = %step.run_id,
1381 step_id = %step.id,
1382 error = %err,
1383 "failed to resolve a waiting signal step"
1384 );
1385 rejected.push(SignalRejected {
1386 run_id: step.run_id,
1387 step_id: step.id,
1388 error: err.to_string(),
1389 });
1390 }
1391 }
1392 }
1393
1394 info!(
1395 signal_id = %stored.id,
1396 signal = %stored.name,
1397 key = %stored.key,
1398 resumed = resumed.len(),
1399 rejected = rejected.len(),
1400 "signal received"
1401 );
1402 self.event_publisher
1403 .publish(Event::SignalReceived(SignalReceivedEvent {
1404 signal_id: stored.id,
1405 name: stored.name.clone(),
1406 key: stored.key.clone(),
1407 resumed_runs: resumed.iter().map(|r| r.run_id).collect(),
1408 at: stored.received_at,
1409 }));
1410
1411 Ok(SignalDelivery {
1412 signal_id: stored.id,
1413 duplicate: false,
1414 resumed,
1415 rejected,
1416 })
1417 }
1418
1419 pub async fn send_signal<S: Signal>(
1457 self: &Arc<Self>,
1458 signal: &S,
1459 key: &str,
1460 idempotency_id: Option<&str>,
1461 ) -> Result<SignalDelivery, EngineError> {
1462 let payload = to_value(signal)?;
1463 self.deliver_signal(NewSignal {
1464 name: S::NAME.to_string(),
1465 key: key.to_string(),
1466 payload,
1467 idempotency_id: idempotency_id.map(str::to_string),
1468 })
1469 .await
1470 }
1471
1472 pub(crate) fn spawn_local_resume(self: &Arc<Self>, run_id: Uuid) {
1478 let engine = Arc::clone(self);
1479 spawn(async move {
1480 if let Err(err) = engine
1481 .store
1482 .update_run_status(run_id, RunStatus::Running)
1483 .await
1484 {
1485 error!(run_id = %run_id, error = %err, "failed to restart a woken run");
1486 return;
1487 }
1488 if let Err(err) = engine.resume_run(run_id).await {
1489 error!(run_id = %run_id, error = %err, "failed to resume a woken run");
1490 }
1491 });
1492 }
1493
1494 pub async fn fail_or_schedule_retry(
1539 &self,
1540 run_id: Uuid,
1541 error: &str,
1542 retryable: bool,
1543 cost_usd: Option<Decimal>,
1544 duration_ms: Option<u64>,
1545 ) -> Result<RunStatus, EngineError> {
1546 let run = self
1547 .store
1548 .get_run(run_id)
1549 .await?
1550 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
1551
1552 let has_attempts_left = run.retry_count < run.max_retries;
1553 let update = if retryable && has_attempts_left {
1554 let backoff = backoff_for_retry(run.retry_count);
1555 let scheduled_at = Utc::now() + TimeDelta::milliseconds(backoff.as_millis() as i64);
1556
1557 info!(
1558 run_id = %run_id,
1559 workflow = %run.workflow_name,
1560 attempt = run.retry_count + 1,
1561 max_retries = run.max_retries,
1562 backoff_secs = backoff.as_secs(),
1563 scheduled_at = %scheduled_at,
1564 "run failed, scheduling retry"
1565 );
1566
1567 RunUpdate {
1568 status: Some(RunStatus::Retrying),
1569 error: Some(error.to_string()),
1570 increment_retry: true,
1571 cost_usd,
1572 duration_ms,
1573 scheduled_at: Some(scheduled_at),
1574 ..RunUpdate::default()
1575 }
1576 } else {
1577 RunUpdate {
1578 status: Some(RunStatus::Failed),
1579 error: Some(error.to_string()),
1580 cost_usd,
1581 duration_ms,
1582 completed_at: Some(Utc::now()),
1583 ..RunUpdate::default()
1584 }
1585 };
1586
1587 let status = update.status.unwrap_or(RunStatus::Failed);
1588 self.store.update_run(run_id, update).await?;
1589 self.fail_orphaned_steps(run_id, error).await?;
1590
1591 Ok(status)
1592 }
1593
1594 pub async fn fail_orphaned_steps(
1608 &self,
1609 run_id: Uuid,
1610 error_message: &str,
1611 ) -> Result<(), EngineError> {
1612 let steps = self.store.list_steps(run_id).await?;
1613 let now = Utc::now();
1614
1615 for step in steps {
1616 if step.status.state.is_terminal() {
1617 continue;
1618 }
1619
1620 let (target_status, error) = match step.status.state {
1621 StepStatus::Running | StepStatus::AwaitingApproval => {
1622 let err = if step.error.is_some() {
1623 None
1624 } else {
1625 Some(error_message.to_string())
1626 };
1627 (StepStatus::Failed, err)
1628 }
1629 StepStatus::Pending => (StepStatus::Skipped, None),
1630 _ => continue,
1631 };
1632
1633 if let Err(e) = self
1634 .store
1635 .update_step(
1636 step.id,
1637 StepUpdate {
1638 status: Some(target_status),
1639 error,
1640 completed_at: Some(now),
1641 ..StepUpdate::default()
1642 },
1643 )
1644 .await
1645 {
1646 warn!(
1647 run_id = %run_id,
1648 step_id = %step.id,
1649 step_name = %step.name,
1650 error = %e,
1651 "failed to cleanup orphaned step"
1652 );
1653 } else {
1654 info!(
1655 run_id = %run_id,
1656 step_id = %step.id,
1657 step_name = %step.name,
1658 from = %step.status.state,
1659 to = %target_status,
1660 "cleaned up orphaned step"
1661 );
1662 }
1663 }
1664
1665 Ok(())
1666 }
1667
1668 async fn release_then_execute(
1673 &self,
1674 run_id: Uuid,
1675 handler: &dyn WorkflowHandler,
1676 ctx: &mut WorkflowContext,
1677 ) -> Result<(), EngineError> {
1678 match self.provider.release_run(&run_id.to_string()).await {
1679 Ok(()) => handler.execute(ctx).await,
1680 Err(e) => Err(EngineError::Operation(OperationError::Agent(e))),
1681 }
1682 }
1683
1684 async fn finalize_run(
1690 &self,
1691 run_id: Uuid,
1692 workflow_name: &str,
1693 result: Result<(), EngineError>,
1694 ctx: &WorkflowContext,
1695 run_start: Instant,
1696 run_labels: HashMap<String, String>,
1697 ) -> Result<WorkflowResult, EngineError> {
1698 let total_duration = ctx.carried_duration_ms() + run_start.elapsed().as_millis() as u64;
1701 let completed_at = Utc::now();
1702
1703 let final_status;
1704 let final_run;
1705
1706 match result {
1707 Ok(()) => {
1708 final_status = if ctx.has_allowed_failure() {
1709 RunStatus::Warning
1710 } else {
1711 RunStatus::Completed
1712 };
1713 final_run = self
1714 .store
1715 .update_run_returning(
1716 run_id,
1717 RunUpdate {
1718 status: Some(final_status),
1719 cost_usd: Some(ctx.total_cost_usd()),
1720 duration_ms: Some(total_duration),
1721 completed_at: Some(completed_at),
1722 ..RunUpdate::default()
1723 },
1724 )
1725 .await?;
1726
1727 info!(
1728 run_id = %run_id,
1729 status = %final_status,
1730 cost_usd = %ctx.total_cost_usd(),
1731 duration_ms = total_duration,
1732 "run completed"
1733 );
1734 }
1735 Err(EngineError::ApprovalRequired {
1736 run_id: approval_run_id,
1737 step_id,
1738 ref message,
1739 }) => {
1740 final_status = RunStatus::AwaitingApproval;
1741 final_run = self
1742 .store
1743 .update_run_returning(
1744 run_id,
1745 RunUpdate {
1746 status: Some(RunStatus::AwaitingApproval),
1747 cost_usd: Some(ctx.total_cost_usd()),
1748 duration_ms: Some(total_duration),
1749 ..RunUpdate::default()
1750 },
1751 )
1752 .await?;
1753
1754 info!(
1755 run_id = %approval_run_id,
1756 step_id = %step_id,
1757 message = %message,
1758 "run awaiting approval"
1759 );
1760
1761 let requirement = self
1763 .store
1764 .get_step(step_id)
1765 .await?
1766 .and_then(|s| s.approval_requirement);
1767 self.event_publisher
1768 .publish(Event::ApprovalRequested(ApprovalRequestedEvent {
1769 run_id: approval_run_id,
1770 step_id,
1771 message: message.clone(),
1772 requirement,
1773 at: Utc::now(),
1774 }));
1775 }
1776 Err(EngineError::HumanInputRequired {
1777 run_id: input_run_id,
1778 step_id,
1779 ref message,
1780 }) => {
1781 final_status = RunStatus::AwaitingApproval;
1782 final_run = self
1783 .store
1784 .update_run_returning(
1785 run_id,
1786 RunUpdate {
1787 status: Some(RunStatus::AwaitingApproval),
1788 cost_usd: Some(ctx.total_cost_usd()),
1789 duration_ms: Some(total_duration),
1790 ..RunUpdate::default()
1791 },
1792 )
1793 .await?;
1794
1795 info!(
1797 run_id = %input_run_id,
1798 step_id = %step_id,
1799 message = %message,
1800 "run awaiting human input"
1801 );
1802 }
1803 Err(EngineError::DelaySleeping {
1804 run_id: delay_run_id,
1805 step_id,
1806 wake_at,
1807 }) => {
1808 final_status = RunStatus::Sleeping;
1809 final_run = self
1810 .store
1811 .update_run_returning(
1812 run_id,
1813 RunUpdate {
1814 status: Some(RunStatus::Sleeping),
1815 cost_usd: Some(ctx.total_cost_usd()),
1816 duration_ms: Some(total_duration),
1817 scheduled_at: Some(wake_at),
1818 ..RunUpdate::default()
1819 },
1820 )
1821 .await?;
1822
1823 info!(
1824 run_id = %delay_run_id,
1825 step_id = %step_id,
1826 wake_at = %wake_at,
1827 "run sleeping until delay elapses"
1828 );
1829 }
1830 Err(EngineError::SignalWaiting {
1831 run_id: wait_run_id,
1832 step_id,
1833 ref step_name,
1834 ref name,
1835 ref key,
1836 deadline_at,
1837 }) => {
1838 final_status = RunStatus::Sleeping;
1839 let waiting = self
1843 .store
1844 .suspend_run_on_signal(run_id, step_id, deadline_at)
1845 .await?;
1846 final_run = self
1847 .store
1848 .update_run_returning(
1849 run_id,
1850 RunUpdate {
1851 cost_usd: Some(ctx.total_cost_usd()),
1852 duration_ms: Some(total_duration),
1853 ..RunUpdate::default()
1854 },
1855 )
1856 .await?;
1857
1858 if waiting {
1859 self.event_publisher
1860 .publish(Event::SignalAwaited(SignalAwaitedEvent {
1861 run_id: wait_run_id,
1862 step_id,
1863 step_name: step_name.clone(),
1864 name: name.clone(),
1865 key: key.clone(),
1866 deadline_at,
1867 at: Utc::now(),
1868 }));
1869 }
1870
1871 info!(
1872 run_id = %wait_run_id,
1873 step_id = %step_id,
1874 signal = %name,
1875 key = %key,
1876 deadline_at = %deadline_at,
1877 waiting,
1878 "run sleeping until a signal arrives"
1879 );
1880 }
1881 Err(err) => {
1882 let guardrail_stop = matches!(
1886 err,
1887 EngineError::RunBudgetExceeded { .. } | EngineError::WorkflowGuardRejected(_)
1888 );
1889
1890 final_status = if guardrail_stop {
1891 if let Err(store_err) = self
1892 .store
1893 .update_run(
1894 run_id,
1895 RunUpdate {
1896 status: Some(RunStatus::Cancelled),
1897 error: Some(err.to_string()),
1898 cost_usd: Some(ctx.total_cost_usd()),
1899 duration_ms: Some(total_duration),
1900 completed_at: Some(completed_at),
1901 ..RunUpdate::default()
1902 },
1903 )
1904 .await
1905 {
1906 error!(run_id = %run_id, store_error = %store_err, "failed to persist run cancellation");
1907 }
1908 if let Err(cleanup_err) = self
1909 .fail_orphaned_steps(run_id, "run stopped: guardrail limit reached")
1910 .await
1911 {
1912 error!(run_id = %run_id, store_error = %cleanup_err, "failed to cleanup orphaned steps");
1913 }
1914 RunStatus::Cancelled
1915 } else {
1916 self.fail_or_schedule_retry(
1917 run_id,
1918 &err.to_string(),
1919 is_run_retryable(&err),
1920 Some(ctx.total_cost_usd()),
1921 Some(total_duration),
1922 )
1923 .await
1924 .unwrap_or_else(|store_err| {
1925 error!(run_id = %run_id, store_error = %store_err, "failed to persist run failure");
1926 RunStatus::Failed
1927 })
1928 };
1929
1930 if matches!(err, EngineError::RunBudgetExceeded { .. }) {
1931 self.on_run_budget_exceeded(workflow_name, run_id, &err);
1932 }
1933
1934 error!(run_id = %run_id, status = %final_status, error = %err, "run stopped");
1935
1936 self.publish_run_status_changed(
1937 workflow_name,
1938 run_id,
1939 final_status,
1940 Some(err.to_string()),
1941 ctx,
1942 total_duration,
1943 run_labels,
1944 );
1945
1946 #[cfg(feature = "prometheus")]
1947 self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
1948
1949 return Err(err);
1950 }
1951 }
1952
1953 self.publish_run_status_changed(
1954 workflow_name,
1955 run_id,
1956 final_status,
1957 None,
1958 ctx,
1959 total_duration,
1960 run_labels,
1961 );
1962
1963 #[cfg(feature = "prometheus")]
1964 self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
1965
1966 Ok(WorkflowResult {
1967 run: final_run,
1968 steps: ctx.step_results().to_vec(),
1969 })
1970 }
1971
1972 #[cfg(feature = "prometheus")]
1974 fn emit_run_metrics(
1975 &self,
1976 workflow_name: &str,
1977 status: RunStatus,
1978 duration_ms: u64,
1979 ctx: &WorkflowContext,
1980 ) {
1981 let status_str = status.to_string();
1982 let wf = workflow_name.to_string();
1983
1984 counter!(RUNS_TOTAL, "workflow" => wf.clone(), "status" => status_str.clone()).increment(1);
1985 histogram!(RUN_DURATION_SECONDS, "workflow" => wf.clone(), "status" => status_str)
1986 .record(duration_ms as f64 / 1000.0);
1987 histogram!(RUN_COST_USD, "workflow" => wf.clone()).record(
1988 ctx.total_cost_usd()
1989 .to_string()
1990 .parse::<f64>()
1991 .unwrap_or(0.0),
1992 );
1993 gauge!(RUNS_ACTIVE, "workflow" => wf).decrement(1.0);
1994 }
1995
1996 fn on_run_budget_exceeded(&self, workflow_name: &str, run_id: Uuid, err: &EngineError) {
2002 let EngineError::RunBudgetExceeded {
2003 limit_usd,
2004 spent_usd,
2005 step_budget_usd,
2006 ..
2007 } = err
2008 else {
2009 return;
2010 };
2011
2012 #[cfg(feature = "prometheus")]
2013 counter!(
2014 RUN_BUDGET_EXCEEDED_TOTAL,
2015 "workflow" => workflow_name.to_string(),
2016 "scope" => "run",
2017 )
2018 .increment(1);
2019
2020 self.event_publisher
2021 .publish(Event::RunBudgetExceeded(RunBudgetExceededEvent {
2022 run_id,
2023 workflow_name: workflow_name.to_string(),
2024 limit_usd: *limit_usd,
2025 spent_usd: *spent_usd,
2026 step_budget_usd: *step_budget_usd,
2027 at: Utc::now(),
2028 }));
2029 }
2030
2031 #[allow(clippy::too_many_arguments)]
2036 fn publish_run_status_changed(
2037 &self,
2038 workflow_name: &str,
2039 run_id: Uuid,
2040 to: RunStatus,
2041 error: Option<String>,
2042 ctx: &WorkflowContext,
2043 duration_ms: u64,
2044 labels: HashMap<String, String>,
2045 ) {
2046 let now = Utc::now();
2047 let cost_usd = ctx.total_cost_usd();
2048 let wf = workflow_name.to_string();
2049
2050 self.event_publisher
2051 .publish(Event::RunStatusChanged(RunStatusChangedEvent {
2052 run_id,
2053 workflow_name: wf.clone(),
2054 from: RunStatus::Running,
2055 to,
2056 error: error.clone(),
2057 cost_usd,
2058 duration_ms,
2059 labels: labels.clone(),
2060 at: now,
2061 }));
2062
2063 if to == RunStatus::Failed {
2064 self.event_publisher
2065 .publish(Event::RunFailed(RunFailedEvent {
2066 run_id,
2067 workflow_name: wf,
2068 error,
2069 cost_usd,
2070 duration_ms,
2071 labels,
2072 at: now,
2073 }));
2074 }
2075 }
2076}
2077
2078impl fmt::Debug for Engine {
2079 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
2080 f.debug_struct("Engine")
2081 .field("handlers", &self.handlers.keys().collect::<Vec<_>>())
2082 .finish_non_exhaustive()
2083 }
2084}
2085
2086#[cfg(test)]
2087mod tests {
2088 use super::*;
2089 use crate::config::ShellConfig;
2090 use crate::handler::{HandlerFuture, WorkflowHandler};
2091 use ironflow_core::providers::claude::ClaudeCodeProvider;
2092 use ironflow_core::providers::record_replay::RecordReplayProvider;
2093 use ironflow_store::memory::InMemoryStore;
2094 use ironflow_store::models::StepStatus;
2095 use serde_json::json;
2096
2097 struct EchoWorkflow;
2099
2100 impl WorkflowHandler for EchoWorkflow {
2101 fn name(&self) -> &str {
2102 "echo-workflow"
2103 }
2104
2105 fn describe(&self) -> WorkflowInfo {
2106 WorkflowInfo {
2107 description: "A simple workflow that echoes hello".to_string(),
2108 source_code: None,
2109 sub_workflows: Vec::new(),
2110 category: None,
2111 version: self.version().map(str::to_string),
2112 compatible_versions: Vec::new(),
2113 input_schema: None,
2114 default_labels: HashMap::new(),
2115 schedule: self.schedule().cloned(),
2116 default_max_cost_usd: self.default_max_cost_usd(),
2117 }
2118 }
2119
2120 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2121 Box::pin(async move {
2122 ctx.shell("greet", ShellConfig::new("echo hello")).await?;
2123 Ok(())
2124 })
2125 }
2126 }
2127
2128 struct FailingWorkflow;
2130
2131 impl WorkflowHandler for FailingWorkflow {
2132 fn name(&self) -> &str {
2133 "failing-workflow"
2134 }
2135
2136 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2137 Box::pin(async move {
2138 ctx.shell("fail", ShellConfig::new("exit 1")).await?;
2139 Ok(())
2140 })
2141 }
2142 }
2143
2144 fn create_test_engine() -> Engine {
2145 let store = Arc::new(InMemoryStore::new());
2146 let inner = ClaudeCodeProvider::new();
2147 let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
2148 inner,
2149 "/tmp/ironflow-fixtures",
2150 ));
2151 Engine::new(store, provider)
2152 }
2153
2154 #[test]
2155 fn engine_new_creates_instance() {
2156 let engine = create_test_engine();
2157 assert_eq!(engine.handler_names().len(), 0);
2158 }
2159
2160 #[test]
2161 fn execution_mode_defaults_to_local() {
2162 let engine = create_test_engine();
2163 assert_eq!(engine.execution_mode(), ExecutionMode::Local);
2164 }
2165
2166 #[test]
2167 fn with_execution_mode_overrides_the_default() {
2168 let engine = create_test_engine().with_execution_mode(ExecutionMode::Workers);
2169 assert_eq!(engine.execution_mode(), ExecutionMode::Workers);
2170 }
2171
2172 #[test]
2173 fn engine_register_handler() {
2174 let mut engine = create_test_engine();
2175 let result = engine.register(EchoWorkflow);
2176 assert!(result.is_ok());
2177 assert_eq!(engine.handler_names().len(), 1);
2178 assert!(engine.handler_names().contains(&"echo-workflow"));
2179 }
2180
2181 #[test]
2182 fn engine_register_duplicate_returns_error() {
2183 let mut engine = create_test_engine();
2184 engine.register(EchoWorkflow).unwrap();
2185 let result = engine.register(EchoWorkflow);
2186 assert!(result.is_err());
2187 }
2188
2189 #[test]
2190 fn engine_get_handler_found() {
2191 let mut engine = create_test_engine();
2192 engine.register(EchoWorkflow).unwrap();
2193 let handler = engine.get_handler("echo-workflow");
2194 assert!(handler.is_some());
2195 }
2196
2197 #[test]
2198 fn engine_get_handler_not_found() {
2199 let engine = create_test_engine();
2200 let handler = engine.get_handler("nonexistent");
2201 assert!(handler.is_none());
2202 }
2203
2204 #[test]
2205 fn engine_handler_names_lists_all() {
2206 let mut engine = create_test_engine();
2207 engine.register(EchoWorkflow).unwrap();
2208 engine.register(FailingWorkflow).unwrap();
2209 let names = engine.handler_names();
2210 assert_eq!(names.len(), 2);
2211 assert!(names.contains(&"echo-workflow"));
2212 assert!(names.contains(&"failing-workflow"));
2213 }
2214
2215 #[test]
2216 fn engine_handler_info_returns_description() {
2217 let mut engine = create_test_engine();
2218 engine.register(EchoWorkflow).unwrap();
2219 let info = engine.handler_info("echo-workflow");
2220 assert!(info.is_some());
2221 let info = info.unwrap();
2222 assert_eq!(info.description, "A simple workflow that echoes hello");
2223 }
2224
2225 struct CategorizedWorkflow;
2226
2227 impl WorkflowHandler for CategorizedWorkflow {
2228 fn name(&self) -> &str {
2229 "categorized"
2230 }
2231 fn category(&self) -> Option<&str> {
2232 Some("data/etl")
2233 }
2234 fn execute<'a>(
2235 &'a self,
2236 _ctx: &'a mut WorkflowContext,
2237 ) -> crate::handler::HandlerFuture<'a> {
2238 Box::pin(async move { Ok(()) })
2239 }
2240 }
2241
2242 #[test]
2243 fn engine_default_describe_propagates_category() {
2244 let mut engine = create_test_engine();
2245 engine.register(CategorizedWorkflow).unwrap();
2246 let info = engine.handler_info("categorized").unwrap();
2247 assert_eq!(info.category.as_deref(), Some("data/etl"));
2248 }
2249
2250 #[test]
2251 fn engine_default_describe_without_category() {
2252 let mut engine = create_test_engine();
2253 engine.register(EchoWorkflow).unwrap();
2254 let info = engine.handler_info("echo-workflow").unwrap();
2255 assert!(info.category.is_none());
2256 }
2257
2258 struct ScheduledWorkflow {
2263 schedule: CronSchedule,
2264 }
2265
2266 impl ScheduledWorkflow {
2267 fn new() -> Self {
2268 Self {
2269 schedule: CronSchedule::new("0 0 * * * *").unwrap(),
2270 }
2271 }
2272 }
2273
2274 impl WorkflowHandler for ScheduledWorkflow {
2275 fn name(&self) -> &str {
2276 "scheduled"
2277 }
2278 fn schedule(&self) -> Option<&CronSchedule> {
2279 Some(&self.schedule)
2280 }
2281 fn execute<'a>(
2282 &'a self,
2283 _ctx: &'a mut WorkflowContext,
2284 ) -> crate::handler::HandlerFuture<'a> {
2285 Box::pin(async move { Ok(()) })
2286 }
2287 }
2288
2289 #[test]
2290 fn engine_default_describe_propagates_schedule() {
2291 let mut engine = create_test_engine();
2292 engine.register(ScheduledWorkflow::new()).unwrap();
2293 let info = engine.handler_info("scheduled").unwrap();
2294 assert_eq!(
2295 info.schedule.as_ref().map(|s| s.as_str()),
2296 Some("0 0 * * * *")
2297 );
2298 }
2299
2300 #[test]
2301 fn engine_default_describe_without_schedule() {
2302 let mut engine = create_test_engine();
2303 engine.register(EchoWorkflow).unwrap();
2304 let info = engine.handler_info("echo-workflow").unwrap();
2305 assert!(info.schedule.is_none());
2306 }
2307
2308 #[test]
2309 fn scheduled_handlers_returns_only_scheduled() {
2310 let mut engine = create_test_engine();
2311 engine.register(EchoWorkflow).unwrap();
2312 engine.register(ScheduledWorkflow::new()).unwrap();
2313 engine.register(FailingWorkflow).unwrap();
2314
2315 let scheduled = engine.scheduled_handlers();
2316 assert_eq!(scheduled.len(), 1);
2317 assert_eq!(scheduled[0].0, "scheduled");
2318 assert_eq!(scheduled[0].1.as_str(), "0 0 * * * *");
2319 }
2320
2321 #[test]
2322 fn scheduled_handlers_empty_when_none_scheduled() {
2323 let mut engine = create_test_engine();
2324 engine.register(EchoWorkflow).unwrap();
2325 engine.register(FailingWorkflow).unwrap();
2326
2327 let scheduled = engine.scheduled_handlers();
2328 assert!(scheduled.is_empty());
2329 }
2330
2331 struct BadCategoryWorkflow(&'static str);
2332
2333 impl WorkflowHandler for BadCategoryWorkflow {
2334 fn name(&self) -> &str {
2335 "bad-category"
2336 }
2337 fn category(&self) -> Option<&str> {
2338 Some(self.0)
2339 }
2340 fn execute<'a>(
2341 &'a self,
2342 _ctx: &'a mut WorkflowContext,
2343 ) -> crate::handler::HandlerFuture<'a> {
2344 Box::pin(async move { Ok(()) })
2345 }
2346 }
2347
2348 #[test]
2349 fn engine_register_rejects_empty_category() {
2350 let mut engine = create_test_engine();
2351 let err = engine.register(BadCategoryWorkflow("")).unwrap_err();
2352 match err {
2353 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty category")),
2354 other => panic!("expected InvalidWorkflow, got {other:?}"),
2355 }
2356 }
2357
2358 #[test]
2359 fn engine_register_rejects_leading_slash_category() {
2360 let mut engine = create_test_engine();
2361 let err = engine
2362 .register(BadCategoryWorkflow("/data/etl"))
2363 .unwrap_err();
2364 match err {
2365 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("leading '/'")),
2366 other => panic!("expected InvalidWorkflow, got {other:?}"),
2367 }
2368 }
2369
2370 #[test]
2371 fn engine_register_rejects_trailing_slash_category() {
2372 let mut engine = create_test_engine();
2373 let err = engine
2374 .register(BadCategoryWorkflow("data/etl/"))
2375 .unwrap_err();
2376 match err {
2377 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("trailing '/'")),
2378 other => panic!("expected InvalidWorkflow, got {other:?}"),
2379 }
2380 }
2381
2382 #[test]
2383 fn engine_register_rejects_double_slash_category() {
2384 let mut engine = create_test_engine();
2385 let err = engine
2386 .register(BadCategoryWorkflow("data//etl"))
2387 .unwrap_err();
2388 match err {
2389 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty segment")),
2390 other => panic!("expected InvalidWorkflow, got {other:?}"),
2391 }
2392 }
2393
2394 #[test]
2395 fn engine_register_rejects_whitespace_only_segment_category() {
2396 let mut engine = create_test_engine();
2397 let err = engine
2398 .register(BadCategoryWorkflow("data/ /etl"))
2399 .unwrap_err();
2400 match err {
2401 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("whitespace-only segment")),
2402 other => panic!("expected InvalidWorkflow, got {other:?}"),
2403 }
2404 }
2405
2406 #[test]
2407 fn engine_register_accepts_valid_nested_category() {
2408 let mut engine = create_test_engine();
2409 assert!(engine.register(CategorizedWorkflow).is_ok());
2410 }
2411
2412 #[tokio::test]
2413 async fn engine_unknown_workflow_returns_error() {
2414 let engine = create_test_engine();
2415 let result = engine
2416 .run_handler("unknown", TriggerKind::Manual, json!({}))
2417 .await;
2418 assert!(result.is_err());
2419 match result {
2420 Err(EngineError::InvalidWorkflow(msg)) => {
2421 assert!(msg.contains("no handler registered"));
2422 }
2423 _ => panic!("expected InvalidWorkflow error"),
2424 }
2425 }
2426
2427 #[tokio::test]
2428 async fn engine_enqueue_handler_creates_pending_run() {
2429 let mut engine = create_test_engine();
2430 engine.register(EchoWorkflow).unwrap();
2431
2432 let run = engine
2433 .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
2434 .await
2435 .unwrap();
2436 assert_eq!(run.status.state, RunStatus::Pending);
2437 assert_eq!(run.workflow_name, "echo-workflow");
2438 }
2439
2440 #[tokio::test]
2441 async fn enqueue_handler_leaves_the_run_unattributed() {
2442 let mut engine = create_test_engine();
2443 engine.register(EchoWorkflow).unwrap();
2444
2445 let run = engine
2446 .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
2447 .await
2448 .unwrap();
2449
2450 assert!(run.created_by.is_none());
2451 }
2452
2453 #[tokio::test]
2454 async fn enqueue_handler_with_options_records_the_author() {
2455 let mut engine = create_test_engine();
2456 engine.register(EchoWorkflow).unwrap();
2457 let actor = RunActor::User {
2458 user_id: Uuid::now_v7(),
2459 };
2460
2461 let run = engine
2462 .enqueue_handler_with_options(
2463 "echo-workflow",
2464 TriggerKind::Api,
2465 json!({}),
2466 EnqueueOptions {
2467 created_by: Some(actor.clone()),
2468 ..Default::default()
2469 },
2470 )
2471 .await
2472 .unwrap()
2473 .into_run();
2474
2475 assert_eq!(run.created_by, Some(actor));
2476 }
2477
2478 #[tokio::test]
2479 async fn enqueue_handler_with_options_accepts_no_author() {
2480 let mut engine = create_test_engine();
2481 engine.register(EchoWorkflow).unwrap();
2482
2483 let run = engine
2484 .enqueue_handler_with_options(
2485 "echo-workflow",
2486 TriggerKind::Cron {
2487 schedule: "0 * * * * *".to_string(),
2488 },
2489 json!({}),
2490 EnqueueOptions::default(),
2491 )
2492 .await
2493 .unwrap()
2494 .into_run();
2495
2496 assert!(run.created_by.is_none());
2497 }
2498
2499 #[tokio::test]
2500 async fn run_handler_leaves_the_run_unattributed() {
2501 let mut engine = create_test_engine();
2502 engine.register(EchoWorkflow).unwrap();
2503
2504 let run = engine
2505 .run_handler("echo-workflow", TriggerKind::Manual, json!({}))
2506 .await
2507 .unwrap()
2508 .run;
2509
2510 assert!(run.created_by.is_none());
2511 }
2512
2513 #[tokio::test]
2514 async fn engine_register_boxed() {
2515 let mut engine = create_test_engine();
2516 let handler: Box<dyn WorkflowHandler> = Box::new(EchoWorkflow);
2517 let result = engine.register_boxed(handler);
2518 assert!(result.is_ok());
2519 assert_eq!(engine.handler_names().len(), 1);
2520 }
2521
2522 #[tokio::test]
2523 async fn engine_store_and_provider_accessors() {
2524 let store = Arc::new(InMemoryStore::new());
2525 let inner = ClaudeCodeProvider::new();
2526 let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
2527 inner,
2528 "/tmp/ironflow-fixtures",
2529 ));
2530 let engine = Engine::new(store.clone(), provider.clone());
2531
2532 let _ = engine.store();
2534 let _ = engine.provider();
2535 }
2536
2537 use crate::operation::{Operation, OperationContext};
2542 use async_trait::async_trait;
2543 use ironflow_core::error::OperationError;
2544 use ironflow_store::models::StepKind;
2545
2546 struct FakeGitlabOp {
2547 project_id: u64,
2548 title: String,
2549 }
2550
2551 #[async_trait]
2552 impl Operation for FakeGitlabOp {
2553 fn kind(&self) -> &str {
2554 "gitlab"
2555 }
2556
2557 async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
2558 Ok(json!({
2559 "issue_id": 42,
2560 "project_id": self.project_id,
2561 "title": self.title,
2562 }))
2563 }
2564
2565 fn input(&self) -> Option<Value> {
2566 Some(json!({
2567 "project_id": self.project_id,
2568 "title": self.title,
2569 }))
2570 }
2571 }
2572
2573 struct FailingOp;
2574
2575 #[async_trait]
2576 impl Operation for FailingOp {
2577 fn kind(&self) -> &str {
2578 "broken-service"
2579 }
2580
2581 async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
2582 Err(OperationError::Http {
2583 status: None,
2584 message: "service unavailable".to_string(),
2585 })
2586 }
2587 }
2588
2589 struct OperationWorkflow;
2590
2591 impl WorkflowHandler for OperationWorkflow {
2592 fn name(&self) -> &str {
2593 "operation-workflow"
2594 }
2595
2596 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2597 Box::pin(async move {
2598 let op = FakeGitlabOp {
2599 project_id: 123,
2600 title: "Bug report".to_string(),
2601 };
2602 ctx.operation("create-issue", &op).await?;
2603 Ok(())
2604 })
2605 }
2606 }
2607
2608 struct FailingOperationWorkflow;
2609
2610 impl WorkflowHandler for FailingOperationWorkflow {
2611 fn name(&self) -> &str {
2612 "failing-operation-workflow"
2613 }
2614
2615 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2616 Box::pin(async move {
2617 ctx.operation("broken-call", &FailingOp).await?;
2618 Ok(())
2619 })
2620 }
2621 }
2622
2623 struct MixedWorkflow;
2624
2625 impl WorkflowHandler for MixedWorkflow {
2626 fn name(&self) -> &str {
2627 "mixed-workflow"
2628 }
2629
2630 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2631 Box::pin(async move {
2632 ctx.shell("build", ShellConfig::new("echo built")).await?;
2633 let op = FakeGitlabOp {
2634 project_id: 456,
2635 title: "Deploy done".to_string(),
2636 };
2637 let result = ctx.operation("notify-gitlab", &op).await?;
2638 assert_eq!(result.output["issue_id"], 42);
2639 Ok(())
2640 })
2641 }
2642 }
2643
2644 #[tokio::test]
2645 async fn operation_step_happy_path() {
2646 let mut engine = create_test_engine();
2647 engine.register(OperationWorkflow).unwrap();
2648
2649 let run = engine
2650 .run_handler("operation-workflow", TriggerKind::Manual, json!({}))
2651 .await
2652 .unwrap()
2653 .run;
2654
2655 assert_eq!(run.status.state, RunStatus::Completed);
2656
2657 let steps = engine.store().list_steps(run.id).await.unwrap();
2658
2659 assert_eq!(steps.len(), 1);
2660 assert_eq!(steps[0].name, "create-issue");
2661 assert_eq!(steps[0].kind, StepKind::Custom("gitlab".to_string()));
2662 assert_eq!(
2663 steps[0].status.state,
2664 ironflow_store::models::StepStatus::Completed
2665 );
2666
2667 let output = steps[0].output.as_ref().unwrap();
2668 assert_eq!(output["issue_id"], 42);
2669 assert_eq!(output["project_id"], 123);
2670
2671 let input = steps[0].input.as_ref().unwrap();
2672 assert_eq!(input["project_id"], 123);
2673 assert_eq!(input["title"], "Bug report");
2674 }
2675
2676 #[tokio::test]
2677 async fn operation_step_failure_marks_run_failed() {
2678 let mut engine = create_test_engine();
2679 engine.register(FailingOperationWorkflow).unwrap();
2680
2681 let result = engine
2682 .run_handler("failing-operation-workflow", TriggerKind::Manual, json!({}))
2683 .await;
2684
2685 assert!(result.is_err());
2686 }
2687
2688 #[tokio::test]
2689 async fn operation_mixed_with_shell_steps() {
2690 let mut engine = create_test_engine();
2691 engine.register(MixedWorkflow).unwrap();
2692
2693 let run = engine
2694 .run_handler("mixed-workflow", TriggerKind::Manual, json!({}))
2695 .await
2696 .unwrap()
2697 .run;
2698
2699 assert_eq!(run.status.state, RunStatus::Completed);
2700
2701 let steps = engine.store().list_steps(run.id).await.unwrap();
2702
2703 assert_eq!(steps.len(), 2);
2704 assert_eq!(steps[0].kind, StepKind::Shell);
2705 assert_eq!(steps[1].kind, StepKind::Custom("gitlab".to_string()));
2706 assert_eq!(steps[0].position, 0);
2707 assert_eq!(steps[1].position, 1);
2708 }
2709
2710 use crate::config::ApprovalConfig;
2715
2716 struct SingleApprovalWorkflow;
2717
2718 impl WorkflowHandler for SingleApprovalWorkflow {
2719 fn name(&self) -> &str {
2720 "single-approval"
2721 }
2722
2723 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2724 Box::pin(async move {
2725 ctx.shell("build", ShellConfig::new("echo built")).await?;
2726 ctx.approval("gate", ApprovalConfig::new("OK?")).await?;
2727 ctx.shell("deploy", ShellConfig::new("echo deployed"))
2728 .await?;
2729 Ok(())
2730 })
2731 }
2732 }
2733
2734 struct DoubleApprovalWorkflow;
2735
2736 impl WorkflowHandler for DoubleApprovalWorkflow {
2737 fn name(&self) -> &str {
2738 "double-approval"
2739 }
2740
2741 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2742 Box::pin(async move {
2743 ctx.shell("build", ShellConfig::new("echo built")).await?;
2744 ctx.approval("staging-gate", ApprovalConfig::new("Deploy staging?"))
2745 .await?;
2746 ctx.shell("deploy-staging", ShellConfig::new("echo staging"))
2747 .await?;
2748 ctx.approval("prod-gate", ApprovalConfig::new("Deploy prod?"))
2749 .await?;
2750 ctx.shell("deploy-prod", ShellConfig::new("echo prod"))
2751 .await?;
2752 Ok(())
2753 })
2754 }
2755 }
2756
2757 #[tokio::test]
2758 async fn approval_pauses_run() {
2759 let mut engine = create_test_engine();
2760 engine.register(SingleApprovalWorkflow).unwrap();
2761
2762 let run = engine
2763 .run_handler("single-approval", TriggerKind::Manual, json!({}))
2764 .await
2765 .unwrap()
2766 .run;
2767
2768 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
2769
2770 let steps = engine.store().list_steps(run.id).await.unwrap();
2771 assert_eq!(steps.len(), 2); assert_eq!(steps[0].kind, StepKind::Shell);
2773 assert_eq!(steps[0].status.state, StepStatus::Completed);
2774 assert_eq!(steps[1].kind, StepKind::Approval);
2775 assert_eq!(steps[1].status.state, StepStatus::AwaitingApproval);
2776 }
2777
2778 #[tokio::test]
2779 async fn approval_resume_completes_run() {
2780 let mut engine = create_test_engine();
2781 engine.register(SingleApprovalWorkflow).unwrap();
2782
2783 let run = engine
2785 .run_handler("single-approval", TriggerKind::Manual, json!({}))
2786 .await
2787 .unwrap()
2788 .run;
2789 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
2790
2791 engine
2793 .store()
2794 .update_run_status(run.id, RunStatus::Running)
2795 .await
2796 .unwrap();
2797
2798 let resumed = engine.resume_run(run.id).await.unwrap().run;
2800 assert_eq!(resumed.status.state, RunStatus::Completed);
2801
2802 let steps = engine.store().list_steps(run.id).await.unwrap();
2803 assert_eq!(steps.len(), 3); assert_eq!(steps[0].name, "build");
2805 assert_eq!(steps[0].status.state, StepStatus::Completed);
2806 assert_eq!(steps[1].name, "gate");
2807 assert_eq!(steps[1].kind, StepKind::Approval);
2808 assert_eq!(steps[1].status.state, StepStatus::Completed);
2809 assert_eq!(steps[2].name, "deploy");
2810 assert_eq!(steps[2].status.state, StepStatus::Completed);
2811 }
2812
2813 #[tokio::test]
2814 async fn double_approval_two_resumes() {
2815 let mut engine = create_test_engine();
2816 engine.register(DoubleApprovalWorkflow).unwrap();
2817
2818 let run = engine
2820 .run_handler("double-approval", TriggerKind::Manual, json!({}))
2821 .await
2822 .unwrap()
2823 .run;
2824 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
2825
2826 let steps = engine.store().list_steps(run.id).await.unwrap();
2827 assert_eq!(steps.len(), 2); engine
2831 .store()
2832 .update_run_status(run.id, RunStatus::Running)
2833 .await
2834 .unwrap();
2835
2836 let resumed = engine.resume_run(run.id).await.unwrap().run;
2837 assert_eq!(resumed.status.state, RunStatus::AwaitingApproval);
2838
2839 let steps = engine.store().list_steps(run.id).await.unwrap();
2840 assert_eq!(steps.len(), 4); engine
2844 .store()
2845 .update_run_status(run.id, RunStatus::Running)
2846 .await
2847 .unwrap();
2848
2849 let final_run = engine.resume_run(run.id).await.unwrap().run;
2850 assert_eq!(final_run.status.state, RunStatus::Completed);
2851
2852 let steps = engine.store().list_steps(run.id).await.unwrap();
2853 assert_eq!(steps.len(), 5);
2854 assert_eq!(steps[0].name, "build");
2855 assert_eq!(steps[1].name, "staging-gate");
2856 assert_eq!(steps[2].name, "deploy-staging");
2857 assert_eq!(steps[3].name, "prod-gate");
2858 assert_eq!(steps[4].name, "deploy-prod");
2859
2860 for step in &steps {
2861 assert_eq!(step.status.state, StepStatus::Completed);
2862 }
2863 }
2864
2865 use ironflow_store::models::{NewStep, StepUpdate, step_trace_id};
2870
2871 async fn create_step_with_status(
2872 store: &Arc<dyn Store>,
2873 run_id: Uuid,
2874 name: &str,
2875 position: u32,
2876 status: StepStatus,
2877 ) -> ironflow_store::models::Step {
2878 let step = store
2879 .create_step(NewStep {
2880 run_id,
2881 trace_id: step_trace_id(run_id, name, position),
2882 name: name.to_string(),
2883 kind: StepKind::Shell,
2884 position,
2885 input: None,
2886 is_error_handler: false,
2887 })
2888 .await
2889 .unwrap();
2890
2891 match status {
2892 StepStatus::Pending => {}
2893 StepStatus::Running => {
2894 store
2895 .update_step(
2896 step.id,
2897 StepUpdate {
2898 status: Some(StepStatus::Running),
2899 ..StepUpdate::default()
2900 },
2901 )
2902 .await
2903 .unwrap();
2904 }
2905 StepStatus::Completed => {
2906 store
2907 .update_step(
2908 step.id,
2909 StepUpdate {
2910 status: Some(StepStatus::Running),
2911 ..StepUpdate::default()
2912 },
2913 )
2914 .await
2915 .unwrap();
2916 store
2917 .update_step(
2918 step.id,
2919 StepUpdate {
2920 status: Some(StepStatus::Completed),
2921 ..StepUpdate::default()
2922 },
2923 )
2924 .await
2925 .unwrap();
2926 }
2927 StepStatus::AwaitingApproval => {
2928 store
2929 .update_step(
2930 step.id,
2931 StepUpdate {
2932 status: Some(StepStatus::Running),
2933 ..StepUpdate::default()
2934 },
2935 )
2936 .await
2937 .unwrap();
2938 store
2939 .update_step(
2940 step.id,
2941 StepUpdate {
2942 status: Some(StepStatus::AwaitingApproval),
2943 ..StepUpdate::default()
2944 },
2945 )
2946 .await
2947 .unwrap();
2948 }
2949 _ => panic!("unsupported status for test helper: {status}"),
2950 }
2951
2952 store.get_step(step.id).await.unwrap().unwrap()
2953 }
2954
2955 #[tokio::test]
2956 async fn fail_orphaned_steps_marks_running_as_failed() {
2957 let engine = create_test_engine();
2958 let run = engine
2959 .store()
2960 .create_run(NewRun {
2961 created_by: None,
2962 workflow_name: "test".to_string(),
2963 trigger: TriggerKind::Manual,
2964 payload: json!({}),
2965 max_retries: 0,
2966 handler_version: None,
2967 labels: HashMap::new(),
2968 scheduled_at: None,
2969 idempotency_key: None,
2970 max_cost_usd: None,
2971 })
2972 .await
2973 .unwrap()
2974 .into_run();
2975
2976 let step = create_step_with_status(
2977 engine.store(),
2978 run.id,
2979 "running-step",
2980 0,
2981 StepStatus::Running,
2982 )
2983 .await;
2984
2985 engine
2986 .fail_orphaned_steps(run.id, "parent run timed out")
2987 .await
2988 .unwrap();
2989
2990 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
2991 assert_eq!(updated.status.state, StepStatus::Failed);
2992 assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
2993 assert!(updated.completed_at.is_some());
2994 }
2995
2996 #[tokio::test]
2997 async fn fail_orphaned_steps_marks_pending_as_skipped() {
2998 let engine = create_test_engine();
2999 let run = engine
3000 .store()
3001 .create_run(NewRun {
3002 created_by: None,
3003 workflow_name: "test".to_string(),
3004 trigger: TriggerKind::Manual,
3005 payload: json!({}),
3006 max_retries: 0,
3007 handler_version: None,
3008 labels: HashMap::new(),
3009 scheduled_at: None,
3010 idempotency_key: None,
3011 max_cost_usd: None,
3012 })
3013 .await
3014 .unwrap()
3015 .into_run();
3016
3017 let step = create_step_with_status(
3018 engine.store(),
3019 run.id,
3020 "pending-step",
3021 0,
3022 StepStatus::Pending,
3023 )
3024 .await;
3025
3026 engine
3027 .fail_orphaned_steps(run.id, "parent run timed out")
3028 .await
3029 .unwrap();
3030
3031 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
3032 assert_eq!(updated.status.state, StepStatus::Skipped);
3033 assert!(updated.error.is_none());
3034 assert!(updated.completed_at.is_some());
3035 }
3036
3037 #[tokio::test]
3038 async fn fail_orphaned_steps_marks_awaiting_approval_as_failed() {
3039 let engine = create_test_engine();
3040 let run = engine
3041 .store()
3042 .create_run(NewRun {
3043 created_by: None,
3044 workflow_name: "test".to_string(),
3045 trigger: TriggerKind::Manual,
3046 payload: json!({}),
3047 max_retries: 0,
3048 handler_version: None,
3049 labels: HashMap::new(),
3050 scheduled_at: None,
3051 idempotency_key: None,
3052 max_cost_usd: None,
3053 })
3054 .await
3055 .unwrap()
3056 .into_run();
3057
3058 let step = create_step_with_status(
3059 engine.store(),
3060 run.id,
3061 "approval-step",
3062 0,
3063 StepStatus::AwaitingApproval,
3064 )
3065 .await;
3066
3067 engine
3068 .fail_orphaned_steps(run.id, "parent run timed out")
3069 .await
3070 .unwrap();
3071
3072 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
3073 assert_eq!(updated.status.state, StepStatus::Failed);
3074 assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
3075 assert!(updated.completed_at.is_some());
3076 }
3077
3078 #[tokio::test]
3079 async fn fail_orphaned_steps_skips_terminal_steps() {
3080 let engine = create_test_engine();
3081 let run = engine
3082 .store()
3083 .create_run(NewRun {
3084 created_by: None,
3085 workflow_name: "test".to_string(),
3086 trigger: TriggerKind::Manual,
3087 payload: json!({}),
3088 max_retries: 0,
3089 handler_version: None,
3090 labels: HashMap::new(),
3091 scheduled_at: None,
3092 idempotency_key: None,
3093 max_cost_usd: None,
3094 })
3095 .await
3096 .unwrap()
3097 .into_run();
3098
3099 let completed_step =
3100 create_step_with_status(engine.store(), run.id, "done", 0, StepStatus::Completed).await;
3101 let running_step =
3102 create_step_with_status(engine.store(), run.id, "in-flight", 1, StepStatus::Running)
3103 .await;
3104
3105 engine
3106 .fail_orphaned_steps(run.id, "parent run timed out")
3107 .await
3108 .unwrap();
3109
3110 let completed = engine
3111 .store()
3112 .get_step(completed_step.id)
3113 .await
3114 .unwrap()
3115 .unwrap();
3116 assert_eq!(completed.status.state, StepStatus::Completed);
3117
3118 let failed = engine
3119 .store()
3120 .get_step(running_step.id)
3121 .await
3122 .unwrap()
3123 .unwrap();
3124 assert_eq!(failed.status.state, StepStatus::Failed);
3125 }
3126
3127 #[tokio::test]
3128 async fn fail_orphaned_steps_mixed_states() {
3129 let engine = create_test_engine();
3130 let run = engine
3131 .store()
3132 .create_run(NewRun {
3133 created_by: None,
3134 workflow_name: "test".to_string(),
3135 trigger: TriggerKind::Manual,
3136 payload: json!({}),
3137 max_retries: 0,
3138 handler_version: None,
3139 labels: HashMap::new(),
3140 scheduled_at: None,
3141 idempotency_key: None,
3142 max_cost_usd: None,
3143 })
3144 .await
3145 .unwrap()
3146 .into_run();
3147
3148 let s_completed =
3149 create_step_with_status(engine.store(), run.id, "step-1", 0, StepStatus::Completed)
3150 .await;
3151 let s_running =
3152 create_step_with_status(engine.store(), run.id, "step-2", 1, StepStatus::Running).await;
3153 let s_pending =
3154 create_step_with_status(engine.store(), run.id, "step-3", 2, StepStatus::Pending).await;
3155
3156 engine.fail_orphaned_steps(run.id, "timeout").await.unwrap();
3157
3158 let r_completed = engine
3159 .store()
3160 .get_step(s_completed.id)
3161 .await
3162 .unwrap()
3163 .unwrap();
3164 assert_eq!(r_completed.status.state, StepStatus::Completed);
3165
3166 let r_running = engine
3167 .store()
3168 .get_step(s_running.id)
3169 .await
3170 .unwrap()
3171 .unwrap();
3172 assert_eq!(r_running.status.state, StepStatus::Failed);
3173 assert_eq!(r_running.error.as_deref(), Some("timeout"));
3174
3175 let r_pending = engine
3176 .store()
3177 .get_step(s_pending.id)
3178 .await
3179 .unwrap()
3180 .unwrap();
3181 assert_eq!(r_pending.status.state, StepStatus::Skipped);
3182 assert!(r_pending.error.is_none());
3183 }
3184
3185 #[tokio::test]
3186 async fn fail_orphaned_steps_no_steps_is_noop() {
3187 let engine = create_test_engine();
3188 let run = engine
3189 .store()
3190 .create_run(NewRun {
3191 created_by: None,
3192 workflow_name: "test".to_string(),
3193 trigger: TriggerKind::Manual,
3194 payload: json!({}),
3195 max_retries: 0,
3196 handler_version: None,
3197 labels: HashMap::new(),
3198 scheduled_at: None,
3199 idempotency_key: None,
3200 max_cost_usd: None,
3201 })
3202 .await
3203 .unwrap()
3204 .into_run();
3205
3206 let result = engine.fail_orphaned_steps(run.id, "timeout").await;
3207 assert!(result.is_ok());
3208 }
3209
3210 #[tokio::test]
3211 async fn fail_orphaned_steps_preserves_existing_error() {
3212 let engine = create_test_engine();
3213 let run = engine
3214 .store()
3215 .create_run(NewRun {
3216 created_by: None,
3217 workflow_name: "test".to_string(),
3218 trigger: TriggerKind::Manual,
3219 payload: json!({}),
3220 max_retries: 0,
3221 handler_version: None,
3222 labels: HashMap::new(),
3223 scheduled_at: None,
3224 idempotency_key: None,
3225 max_cost_usd: None,
3226 })
3227 .await
3228 .unwrap()
3229 .into_run();
3230
3231 let step_with_error = create_step_with_status(
3232 engine.store(),
3233 run.id,
3234 "already-errored",
3235 0,
3236 StepStatus::Running,
3237 )
3238 .await;
3239
3240 engine
3241 .store()
3242 .update_step(
3243 step_with_error.id,
3244 StepUpdate {
3245 error: Some("real error from provider".to_string()),
3246 ..StepUpdate::default()
3247 },
3248 )
3249 .await
3250 .unwrap();
3251
3252 let step_no_error = create_step_with_status(
3253 engine.store(),
3254 run.id,
3255 "no-error-yet",
3256 1,
3257 StepStatus::Running,
3258 )
3259 .await;
3260
3261 engine
3262 .fail_orphaned_steps(run.id, "parent run failed")
3263 .await
3264 .unwrap();
3265
3266 let updated_with = engine
3267 .store()
3268 .get_step(step_with_error.id)
3269 .await
3270 .unwrap()
3271 .unwrap();
3272 assert_eq!(updated_with.status.state, StepStatus::Failed);
3273 assert_eq!(
3274 updated_with.error.as_deref(),
3275 Some("real error from provider"),
3276 );
3277
3278 let updated_without = engine
3279 .store()
3280 .get_step(step_no_error.id)
3281 .await
3282 .unwrap()
3283 .unwrap();
3284 assert_eq!(updated_without.status.state, StepStatus::Failed);
3285 assert_eq!(updated_without.error.as_deref(), Some("parent run failed"),);
3286 }
3287}