1use std::collections::{HashMap, HashSet};
10use std::fmt;
11use std::sync::{Arc, Mutex};
12use std::time::Instant;
13
14use chrono::{DateTime, TimeDelta, Utc};
15use rust_decimal::Decimal;
16use serde_json::{Value, to_value};
17use tokio::spawn;
18use tracing::{error, info, warn};
19use uuid::Uuid;
20
21use ironflow_core::error::OperationError;
22#[cfg(feature = "prometheus")]
23use ironflow_core::metric_names::{
24 RUN_BUDGET_EXCEEDED_TOTAL, RUN_COST_USD, RUN_DURATION_SECONDS, RUNS_ACTIVE, RUNS_TOTAL,
25};
26use ironflow_core::provider::{AgentProvider, LABEL_ROOT_RUN_ID};
27use ironflow_store::error::StoreError;
28use ironflow_store::models::{
29 ConcurrencyLimit, LeaseUpdate, NewRun, NewSignal, ProviderKind, Run, RunActor, RunCreation,
30 RunFilter, RunStatus, RunUpdate, SignalInsert, SignalStepResolution, StepStatus, StepUpdate,
31 TriggerKind, normalize_worker_tags, validate_concurrency_limits, validate_priority,
32 validate_worker_tags,
33};
34use ironflow_store::store::Store;
35#[cfg(feature = "prometheus")]
36use metrics::{counter, gauge, histogram};
37
38use crate::artifact::ArtifactSink;
39use crate::budget::{BudgetConfig, month_start};
40use crate::context::{PARENT_RUN_ID_LABEL, WorkflowContext, interrupt_running_steps};
41use crate::error::EngineError;
42use crate::executor::{StepInterceptor, StepResult};
43use crate::guard::{WorkflowGuardConfig, new_shared_guard_state};
44use crate::handler::{WorkflowHandler, WorkflowInfo, clamp_priority};
45use crate::log_sender::LogSender;
46use crate::notify::{
47 ApprovalRequestedEvent, Event, EventPublisher, EventSubscriber, RunBudgetExceededEvent,
48 RunFailedEvent, RunStatusChangedEvent, SignalAwaitedEvent, SignalReceivedEvent,
49 WorkflowEventBus,
50};
51use crate::plan::{
52 ExecutionPlan, PlanOptions, PlanRecorder, SharedPlanRecorder, estimate_durations, lock_plan,
53};
54use crate::retry_policy::{backoff_for_retry, is_run_retryable};
55use crate::schedule::CronSchedule;
56use crate::signal::{
57 Signal, SignalDelivery, SignalRejected, SignalResumed, received_output, validate_step_payload,
58};
59use ironflow_core::decision::DecisionProvider;
60
61#[derive(Debug, Clone)]
80pub struct WorkflowResult {
81 pub run: Run,
83 pub steps: Vec<StepResult>,
85}
86
87#[derive(Debug, Clone, Default)]
106pub struct EnqueueOptions {
107 pub max_retries: u32,
109 pub labels: HashMap<String, String>,
111 pub scheduled_at: Option<DateTime<Utc>>,
114 pub max_cost_usd: Option<Decimal>,
118 pub created_by: Option<RunActor>,
121 pub idempotency_key: Option<String>,
127 pub concurrency_key: Option<String>,
133 pub concurrency_limits: Vec<ConcurrencyLimit>,
141 pub priority: Option<i16>,
150 pub worker_tags: Vec<String>,
156}
157
158#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
169pub enum ExecutionMode {
170 #[default]
174 Local,
175 Workers,
178}
179
180pub struct Engine {
221 store: Arc<dyn Store>,
222 provider: Arc<dyn AgentProvider>,
223 handlers: HashMap<String, Arc<dyn WorkflowHandler>>,
224 event_publisher: EventPublisher,
225 log_sender: Option<LogSender>,
226 budget: BudgetConfig,
227 artifact_sink: Option<Arc<dyn ArtifactSink>>,
228 guard_config: Option<WorkflowGuardConfig>,
229 event_bus: Option<WorkflowEventBus>,
230 decision_provider: Option<Arc<dyn DecisionProvider>>,
231 step_interceptor: Option<Arc<dyn StepInterceptor>>,
232 execution_mode: ExecutionMode,
233 worker_tags: Option<Arc<Vec<String>>>,
234}
235
236fn validate_category(handler_name: &str, category: &str) -> Result<(), EngineError> {
246 let reject = |reason: &str| {
247 Err(EngineError::InvalidWorkflow(format!(
248 "handler '{handler_name}' has invalid category '{category}': {reason}"
249 )))
250 };
251
252 if category.is_empty() {
253 return reject("empty category");
254 }
255 if category.starts_with('/') {
256 return reject("leading '/'");
257 }
258 if category.ends_with('/') {
259 return reject("trailing '/'");
260 }
261 for segment in category.split('/') {
262 if segment.is_empty() {
263 return reject("empty segment (double '/')");
264 }
265 if segment.trim().is_empty() {
266 return reject("whitespace-only segment");
267 }
268 }
269 Ok(())
270}
271
272fn chain_label(run: &Run, key: &str) -> Option<Uuid> {
277 if !matches!(run.trigger, TriggerKind::Workflow) {
278 return None;
279 }
280 let id = Uuid::parse_str(run.labels.get(key)?).ok()?;
281 (id != run.id).then_some(id)
282}
283
284pub fn chain_root(run: &Run) -> Option<Uuid> {
304 chain_label(run, LABEL_ROOT_RUN_ID)
305}
306
307fn chain_parent(run: &Run) -> Option<Uuid> {
309 chain_label(run, PARENT_RUN_ID_LABEL)
310}
311
312impl Engine {
313 pub fn new(store: Arc<dyn Store>, provider: Arc<dyn AgentProvider>) -> Self {
329 Self {
330 store,
331 provider,
332 handlers: HashMap::new(),
333 event_publisher: EventPublisher::new(),
334 log_sender: None,
335 budget: BudgetConfig::new(),
336 artifact_sink: None,
337 guard_config: None,
338 event_bus: None,
339 decision_provider: None,
340 step_interceptor: None,
341 execution_mode: ExecutionMode::default(),
342 worker_tags: None,
343 }
344 }
345
346 pub fn with_decision_provider(mut self, provider: Arc<dyn DecisionProvider>) -> Self {
367 self.decision_provider = Some(provider);
368 self
369 }
370
371 pub fn with_step_interceptor(mut self, interceptor: Arc<dyn StepInterceptor>) -> Self {
395 self.step_interceptor = Some(interceptor);
396 self
397 }
398
399 pub fn step_interceptor(&self) -> Option<&Arc<dyn StepInterceptor>> {
401 self.step_interceptor.as_ref()
402 }
403
404 pub fn with_budget_config(mut self, budget: BudgetConfig) -> Self {
425 self.budget = budget;
426 self
427 }
428
429 pub fn budget_config(&self) -> &BudgetConfig {
431 &self.budget
432 }
433
434 pub fn with_guard_config(mut self, config: WorkflowGuardConfig) -> Self {
456 self.guard_config = Some(config);
457 self
458 }
459
460 pub fn guard_config(&self) -> Option<&WorkflowGuardConfig> {
462 self.guard_config.as_ref()
463 }
464
465 pub fn with_execution_mode(mut self, mode: ExecutionMode) -> Self {
486 self.execution_mode = mode;
487 self
488 }
489
490 pub fn execution_mode(&self) -> ExecutionMode {
492 self.execution_mode
493 }
494
495 pub fn set_worker_tags(&mut self, tags: Vec<String>) {
518 self.worker_tags = Some(Arc::new(tags));
519 }
520
521 pub fn worker_tags(&self) -> Option<&[String]> {
539 self.worker_tags.as_deref().map(Vec::as_slice)
540 }
541
542 pub fn set_log_sender(&mut self, sender: LogSender) {
548 self.log_sender = Some(sender);
549 }
550
551 pub fn set_artifact_sink(&mut self, sink: Arc<dyn ArtifactSink>) {
569 self.artifact_sink = Some(sink);
570 }
571
572 pub fn artifact_sink(&self) -> Option<&Arc<dyn ArtifactSink>> {
574 self.artifact_sink.as_ref()
575 }
576
577 pub fn set_event_bus(&mut self, bus: WorkflowEventBus) {
594 self.event_bus = Some(bus);
595 }
596
597 pub fn event_bus(&self) -> Option<&WorkflowEventBus> {
599 self.event_bus.as_ref()
600 }
601
602 pub fn store(&self) -> &Arc<dyn Store> {
604 &self.store
605 }
606
607 pub fn provider(&self) -> &Arc<dyn AgentProvider> {
609 &self.provider
610 }
611
612 fn build_context(&self, run: &Run) -> WorkflowContext {
621 let handlers = self.handlers.clone();
622 let resolver: crate::context::HandlerResolver =
623 Arc::new(move |name: &str| handlers.get(name).cloned());
624 let mut ctx = WorkflowContext::with_handler_resolver(
625 run.id,
626 run.workflow_name.clone(),
627 self.store.clone(),
628 self.provider.clone(),
629 resolver,
630 );
631 ctx.carry_over_run_totals(run.retry_count + 1, run.cost_usd, run.duration_ms);
632 ctx.set_max_cost_usd(run.max_cost_usd);
633 ctx.set_run_created_at(run.created_at);
634 if let Some(ref sender) = self.log_sender {
635 ctx.set_log_sender(sender.clone());
636 }
637 if let Some(ref sink) = self.artifact_sink {
638 ctx.set_artifact_sink(sink.clone());
639 }
640 if let Some(ref bus) = self.event_bus {
641 ctx.set_event_bus(bus.clone());
642 }
643 if let Some(ref provider) = self.decision_provider {
644 ctx.set_decision_provider(provider.clone());
645 }
646 if let Some(ref interceptor) = self.step_interceptor {
647 ctx.set_step_interceptor(interceptor.clone());
648 }
649 if let Some(ref tags) = self.worker_tags {
650 ctx.set_worker_tags(tags.clone());
651 }
652 ctx
653 }
654
655 fn build_context_with_guard(
661 &self,
662 run: &Run,
663 handler: &dyn WorkflowHandler,
664 ) -> WorkflowContext {
665 let mut ctx = self.build_context(run);
666 let guard_config = handler.guard_config().or_else(|| self.guard_config.clone());
667 if let Some(config) = guard_config {
668 ctx.set_guard(config, new_shared_guard_state());
669 }
670 ctx
671 }
672
673 async fn check_monthly_quota(&self, workflow_name: &str) -> Result<(), EngineError> {
684 let Some(limit) = self.budget.monthly_cost_limit_usd else {
685 return Ok(());
686 };
687
688 let stats = self
689 .store
690 .get_stats(RunFilter {
691 created_after: Some(month_start(Utc::now())),
692 ..RunFilter::default()
693 })
694 .await?;
695
696 if stats.total_cost_usd < limit {
697 return Ok(());
698 }
699
700 warn!(
701 workflow = %workflow_name,
702 limit_usd = %limit,
703 spent_usd = %stats.total_cost_usd,
704 "monthly cost quota exhausted, refusing new run"
705 );
706
707 #[cfg(feature = "prometheus")]
708 counter!(
709 RUN_BUDGET_EXCEEDED_TOTAL,
710 "workflow" => workflow_name.to_string(),
711 "scope" => "monthly",
712 )
713 .increment(1);
714
715 Err(EngineError::MonthlyBudgetExceeded {
716 limit_usd: limit,
717 spent_usd: stats.total_cost_usd,
718 })
719 }
720
721 pub fn register(&mut self, handler: impl WorkflowHandler + 'static) -> Result<(), EngineError> {
765 let name = handler.name().to_string();
766 if self.handlers.contains_key(&name) {
767 return Err(EngineError::InvalidWorkflow(format!(
768 "handler '{}' already registered",
769 name
770 )));
771 }
772 if let Some(category) = handler.category() {
773 validate_category(&name, category)?;
774 }
775 self.handlers.insert(name, Arc::new(handler));
776 Ok(())
777 }
778
779 pub fn register_boxed(&mut self, handler: Box<dyn WorkflowHandler>) -> Result<(), EngineError> {
786 let name = handler.name().to_string();
787 if self.handlers.contains_key(&name) {
788 return Err(EngineError::InvalidWorkflow(format!(
789 "handler '{}' already registered",
790 name
791 )));
792 }
793 if let Some(category) = handler.category() {
794 validate_category(&name, category)?;
795 }
796 self.handlers.insert(name, Arc::from(handler));
797 Ok(())
798 }
799
800 pub fn get_handler(&self, name: &str) -> Option<&Arc<dyn WorkflowHandler>> {
802 self.handlers.get(name)
803 }
804
805 pub fn handler_names(&self) -> Vec<&str> {
807 self.handlers.keys().map(|s| s.as_str()).collect()
808 }
809
810 pub fn handler_info(&self, name: &str) -> Option<WorkflowInfo> {
812 self.handlers.get(name).map(|h| h.describe())
813 }
814
815 pub fn scheduled_handlers(&self) -> Vec<(&str, &CronSchedule)> {
839 self.handlers
840 .iter()
841 .filter_map(|(name, handler)| handler.schedule().map(|sched| (name.as_str(), sched)))
842 .collect()
843 }
844
845 pub fn subscribe(
870 &mut self,
871 subscriber: impl EventSubscriber + 'static,
872 event_types: &[&'static str],
873 ) {
874 self.event_publisher.subscribe(subscriber, event_types);
875 }
876
877 pub fn event_publisher(&self) -> &EventPublisher {
882 &self.event_publisher
883 }
884
885 #[tracing::instrument(name = "engine.run_handler", skip_all, fields(workflow = %handler_name))]
915 pub async fn run_handler(
916 &self,
917 handler_name: &str,
918 trigger: TriggerKind,
919 payload: Value,
920 ) -> Result<WorkflowResult, EngineError> {
921 let handler = self
922 .handlers
923 .get(handler_name)
924 .ok_or_else(|| {
925 EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
926 })?
927 .clone();
928
929 self.check_monthly_quota(handler_name).await?;
930
931 let handler_version = handler.version().map(str::to_string);
932 let max_cost_usd = self
933 .budget
934 .resolve_run_cap(None, handler.default_max_cost_usd());
935 let run = self
936 .store
937 .create_run(NewRun {
938 created_by: None,
939 workflow_name: handler_name.to_string(),
940 trigger,
941 payload,
942 max_retries: 0,
943 handler_version,
944 labels: handler.default_labels(),
945 scheduled_at: None,
946 idempotency_key: None,
947 concurrency_key: None,
948 priority: clamp_priority(handler.priority()),
949 concurrency_limits: Vec::new(),
950 max_cost_usd,
951 worker_tags: normalize_worker_tags(handler.required_worker_tags()),
952 })
953 .await?
954 .into_run();
955
956 let run_id = run.id;
957 info!(run_id = %run_id, handler_version = run.handler_version.as_deref().unwrap_or(""), "run created");
958
959 self.store
960 .update_run_status(run_id, RunStatus::Running)
961 .await?;
962
963 #[cfg(feature = "prometheus")]
964 gauge!(RUNS_ACTIVE, "workflow" => handler_name.to_string()).increment(1.0);
965
966 let run_start = Instant::now();
967 let mut ctx = self.build_context_with_guard(&run, handler.as_ref());
968
969 let result = handler.execute(&mut ctx).await;
970 self.finalize_run(run_id, handler_name, result, &ctx, run_start, run.labels)
971 .await
972 }
973
974 #[tracing::instrument(name = "engine.plan_handler", skip_all, fields(workflow = %handler_name))]
1016 pub async fn plan_handler(
1017 &self,
1018 handler_name: &str,
1019 payload: Value,
1020 options: PlanOptions,
1021 ) -> Result<ExecutionPlan, EngineError> {
1022 if options.max_depth == 0 {
1023 return Err(EngineError::InvalidWorkflow(
1024 "max_depth must be at least 1".to_string(),
1025 ));
1026 }
1027
1028 let handler = self
1029 .handlers
1030 .get(handler_name)
1031 .ok_or_else(|| {
1032 EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
1033 })?
1034 .clone();
1035
1036 let estimates = if options.estimate_durations {
1037 estimate_durations(&self.store, handler_name, options.sample_runs).await?
1038 } else {
1039 HashMap::new()
1040 };
1041
1042 let shared: SharedPlanRecorder = Arc::new(Mutex::new(PlanRecorder::new(
1043 handler_name.to_string(),
1044 payload,
1045 options.max_depth,
1046 estimates,
1047 )));
1048
1049 let handlers = self.handlers.clone();
1052 let resolver: crate::context::HandlerResolver =
1053 Arc::new(move |name: &str| handlers.get(name).cloned());
1054 let mut ctx = WorkflowContext::with_handler_resolver(
1055 Uuid::now_v7(),
1056 handler_name.to_string(),
1057 self.store.clone(),
1058 self.provider.clone(),
1059 resolver,
1060 );
1061 ctx.set_plan(shared.clone());
1062
1063 if let Err(err) = handler.execute(&mut ctx).await {
1064 lock_plan(&shared).fail(err.to_string());
1065 }
1066 drop(ctx);
1067
1068 let plan = match Arc::try_unwrap(shared) {
1069 Ok(mutex) => mutex
1070 .into_inner()
1071 .unwrap_or_else(|poisoned| poisoned.into_inner())
1072 .into_plan(),
1073 Err(shared) => lock_plan(&shared).snapshot(),
1074 };
1075
1076 info!(
1077 workflow = %handler_name,
1078 steps = plan.steps.len(),
1079 truncated = plan.truncated,
1080 "execution plan built"
1081 );
1082
1083 Ok(plan)
1084 }
1085
1086 #[tracing::instrument(name = "engine.enqueue_handler", skip_all, fields(workflow = %handler_name))]
1097 pub async fn enqueue_handler(
1098 &self,
1099 handler_name: &str,
1100 trigger: TriggerKind,
1101 payload: Value,
1102 max_retries: u32,
1103 ) -> Result<Run, EngineError> {
1104 self.enqueue_handler_with_options(
1105 handler_name,
1106 trigger,
1107 payload,
1108 EnqueueOptions {
1109 max_retries,
1110 ..Default::default()
1111 },
1112 )
1113 .await
1114 .map(RunCreation::into_run)
1115 }
1116
1117 #[tracing::instrument(name = "engine.enqueue_handler_with_options", skip_all, fields(workflow = %handler_name))]
1169 pub async fn enqueue_handler_with_options(
1170 &self,
1171 handler_name: &str,
1172 trigger: TriggerKind,
1173 payload: Value,
1174 options: EnqueueOptions,
1175 ) -> Result<RunCreation, EngineError> {
1176 let EnqueueOptions {
1177 max_retries,
1178 labels,
1179 scheduled_at,
1180 max_cost_usd,
1181 created_by,
1182 idempotency_key,
1183 concurrency_key,
1184 concurrency_limits,
1185 priority,
1186 worker_tags,
1187 } = options;
1188
1189 validate_concurrency_limits(&concurrency_limits)
1192 .map_err(EngineError::InvalidConcurrencyLimit)?;
1193 if let Some(priority) = priority {
1194 validate_priority(priority).map_err(EngineError::InvalidPriority)?;
1195 }
1196 validate_worker_tags(&worker_tags).map_err(EngineError::InvalidWorkerTag)?;
1197
1198 let handler = self.handlers.get(handler_name).ok_or_else(|| {
1199 EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
1200 })?;
1201
1202 self.check_monthly_quota(handler_name).await?;
1203
1204 let handler_version = handler.version().map(str::to_string);
1205 let mut merged_labels = handler.default_labels();
1206 merged_labels.extend(labels);
1207 let resolved_cap = self
1208 .budget
1209 .resolve_run_cap(max_cost_usd, handler.default_max_cost_usd());
1210 let priority = priority.unwrap_or_else(|| clamp_priority(handler.priority()));
1211 let required_tags = normalize_worker_tags(
1212 handler
1213 .required_worker_tags()
1214 .into_iter()
1215 .chain(worker_tags),
1216 );
1217
1218 let creation = self
1219 .store
1220 .create_run(NewRun {
1221 workflow_name: handler_name.to_string(),
1222 trigger,
1223 payload,
1224 max_retries,
1225 handler_version,
1226 labels: merged_labels,
1227 scheduled_at,
1228 created_by,
1229 idempotency_key,
1230 concurrency_key,
1231 priority,
1232 concurrency_limits,
1233 max_cost_usd: resolved_cap,
1234 worker_tags: required_tags,
1235 })
1236 .await?;
1237
1238 match &creation {
1239 RunCreation::Created(run) => info!(
1240 run_id = %run.id,
1241 workflow = %handler_name,
1242 max_cost_usd = ?resolved_cap,
1243 "handler run enqueued"
1244 ),
1245 RunCreation::Existing(run) => info!(
1246 run_id = %run.id,
1247 workflow = %handler_name,
1248 "idempotent replay, nothing enqueued"
1249 ),
1250 }
1251
1252 Ok(creation)
1253 }
1254
1255 #[tracing::instrument(name = "engine.execute_handler_run", skip_all, fields(run_id = %run_id))]
1275 pub async fn execute_handler_run(&self, run_id: Uuid) -> Result<WorkflowResult, EngineError> {
1276 let run = self
1277 .store
1278 .get_run(run_id)
1279 .await?
1280 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
1281
1282 if let Some(root_run_id) = chain_root(&run) {
1283 return self.resume_chain(run, root_run_id).await;
1284 }
1285
1286 let handler = self
1287 .handlers
1288 .get(&run.workflow_name)
1289 .ok_or_else(|| {
1290 EngineError::InvalidWorkflow(format!(
1291 "no handler registered: {}",
1292 run.workflow_name
1293 ))
1294 })?
1295 .clone();
1296
1297 #[cfg(feature = "prometheus")]
1298 gauge!(RUNS_ACTIVE, "workflow" => run.workflow_name.clone()).increment(1.0);
1299
1300 let run_start = Instant::now();
1301 let mut ctx = self.build_context_with_guard(&run, handler.as_ref());
1302
1303 let result = if handler.is_version_compatible(run.handler_version.as_deref()) {
1316 ctx.load_replay_steps().await?;
1317 self.release_then_execute(run_id, handler.as_ref(), &mut ctx)
1318 .await
1319 } else {
1320 Err(EngineError::HandlerVersionMismatch {
1321 run_id,
1322 workflow_name: run.workflow_name.clone(),
1323 run_version: run
1324 .handler_version
1325 .clone()
1326 .unwrap_or_else(|| "unknown".to_string()),
1327 current_version: handler
1328 .version()
1329 .map(str::to_string)
1330 .unwrap_or_else(|| "unknown".to_string()),
1331 })
1332 };
1333
1334 self.finalize_run(
1335 run_id,
1336 &run.workflow_name,
1337 result,
1338 &ctx,
1339 run_start,
1340 run.labels,
1341 )
1342 .await
1343 }
1344
1345 #[tracing::instrument(name = "engine.execute_run", skip_all, fields(run_id = %run_id))]
1353 pub async fn execute_run(&self, run_id: Uuid) -> Result<WorkflowResult, EngineError> {
1354 self.execute_handler_run(run_id).await
1355 }
1356
1357 #[tracing::instrument(name = "engine.resume_run", skip_all, fields(run_id = %run_id))]
1386 pub async fn resume_run(&self, run_id: Uuid) -> Result<WorkflowResult, EngineError> {
1387 let run = self
1388 .store
1389 .get_run(run_id)
1390 .await?
1391 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
1392
1393 if let Some(root_run_id) = chain_root(&run) {
1394 return self.resume_chain(run, root_run_id).await;
1395 }
1396
1397 self.resume_loaded_run(run).await
1398 }
1399
1400 async fn resume_chain(
1415 &self,
1416 child: Run,
1417 root_run_id: Uuid,
1418 ) -> Result<WorkflowResult, EngineError> {
1419 let child_run_id = child.id;
1420 let lease = child.worker_id.zip(child.lease_expires_at);
1421 let root = self
1422 .store
1423 .get_run(root_run_id)
1424 .await?
1425 .ok_or(EngineError::Store(StoreError::RunNotFound(root_run_id)))?;
1426
1427 match root.status.state {
1428 RunStatus::AwaitingApproval | RunStatus::Pending => {
1429 self.move_root_to_running(root_run_id, lease.as_ref())
1430 .await?;
1431 }
1432 RunStatus::Sleeping => {
1433 self.store
1434 .update_run_status(root_run_id, RunStatus::Pending)
1435 .await?;
1436 self.move_root_to_running(root_run_id, lease.as_ref())
1437 .await?;
1438 }
1439 other => {
1440 let reason = format!(
1441 "cannot resume child run {child_run_id}: root run {root_run_id} is {other}"
1442 );
1443 if let Err(err) = self
1444 .fail_or_schedule_retry(child_run_id, &reason, false, None, None)
1445 .await
1446 {
1447 error!(
1448 run_id = %child_run_id,
1449 error = %err,
1450 "failed to fail a child run whose root cannot resume"
1451 );
1452 }
1453 return Err(EngineError::InvalidWorkflow(reason));
1454 }
1455 }
1456
1457 if lease.is_some() {
1458 self.store
1459 .update_run(
1460 child_run_id,
1461 RunUpdate {
1462 lease: Some(LeaseUpdate::Release),
1463 ..RunUpdate::default()
1464 },
1465 )
1466 .await?;
1467 }
1468
1469 info!(
1470 run_id = %child_run_id,
1471 root_run_id = %root_run_id,
1472 lease_transferred = lease.is_some(),
1473 "child run resumed through its root run"
1474 );
1475
1476 let root = self
1477 .store
1478 .get_run(root_run_id)
1479 .await?
1480 .ok_or(EngineError::Store(StoreError::RunNotFound(root_run_id)))?;
1481 self.resume_loaded_run(root).await
1482 }
1483
1484 async fn move_root_to_running(
1490 &self,
1491 root_run_id: Uuid,
1492 lease: Option<&(String, DateTime<Utc>)>,
1493 ) -> Result<(), EngineError> {
1494 match lease {
1495 Some((worker_id, expires_at)) => {
1496 self.store
1497 .update_run(
1498 root_run_id,
1499 RunUpdate {
1500 status: Some(RunStatus::Running),
1501 lease: Some(LeaseUpdate::Set {
1502 worker_id: worker_id.clone(),
1503 expires_at: *expires_at,
1504 }),
1505 ..RunUpdate::default()
1506 },
1507 )
1508 .await?;
1509 }
1510 None => {
1511 self.store
1512 .update_run_status(root_run_id, RunStatus::Running)
1513 .await?;
1514 }
1515 }
1516 Ok(())
1517 }
1518
1519 async fn resume_loaded_run(&self, run: Run) -> Result<WorkflowResult, EngineError> {
1521 let run_id = run.id;
1522 let handler = self
1523 .handlers
1524 .get(&run.workflow_name)
1525 .ok_or_else(|| {
1526 EngineError::InvalidWorkflow(format!(
1527 "no handler registered: {}",
1528 run.workflow_name
1529 ))
1530 })?
1531 .clone();
1532
1533 info!(run_id = %run_id, workflow = %run.workflow_name, "resuming run after approval");
1534
1535 let run_start = Instant::now();
1536 let mut ctx = self.build_context_with_guard(&run, handler.as_ref());
1537
1538 let result = if handler.is_version_compatible(run.handler_version.as_deref()) {
1539 ctx.load_replay_steps().await?;
1540 self.release_then_execute(run_id, handler.as_ref(), &mut ctx)
1541 .await
1542 } else {
1543 Err(EngineError::HandlerVersionMismatch {
1544 run_id,
1545 workflow_name: run.workflow_name.clone(),
1546 run_version: run
1547 .handler_version
1548 .clone()
1549 .unwrap_or_else(|| "unknown".to_string()),
1550 current_version: handler
1551 .version()
1552 .map(str::to_string)
1553 .unwrap_or_else(|| "unknown".to_string()),
1554 })
1555 };
1556
1557 self.finalize_run(
1558 run_id,
1559 &run.workflow_name,
1560 result,
1561 &ctx,
1562 run_start,
1563 run.labels,
1564 )
1565 .await
1566 }
1567
1568 pub async fn deliver_signal(
1614 self: &Arc<Self>,
1615 signal: NewSignal,
1616 ) -> Result<SignalDelivery, EngineError> {
1617 if signal.name.trim().is_empty() {
1618 return Err(EngineError::InvalidSignal(
1619 "signal name must not be empty".to_string(),
1620 ));
1621 }
1622 if signal.key.trim().is_empty() {
1623 return Err(EngineError::InvalidSignal(
1624 "signal key must not be empty".to_string(),
1625 ));
1626 }
1627
1628 let stored = match self.store.insert_signal(signal).await? {
1629 SignalInsert::Created(stored) => stored,
1630 SignalInsert::Duplicate(existing) => {
1631 info!(
1632 signal_id = %existing.id,
1633 signal = %existing.name,
1634 key = %existing.key,
1635 "duplicate signal ignored"
1636 );
1637 return Ok(SignalDelivery {
1638 signal_id: existing.id,
1639 duplicate: true,
1640 resumed: Vec::new(),
1641 rejected: Vec::new(),
1642 });
1643 }
1644 };
1645
1646 let waiters = self
1647 .store
1648 .list_signal_waiters(&stored.name, &stored.key)
1649 .await?;
1650 let mut resumed = Vec::new();
1651 let mut rejected = Vec::new();
1652
1653 for step in waiters {
1654 if let Err(error) = validate_step_payload(step.input.as_ref(), &stored.payload) {
1655 rejected.push(SignalRejected {
1656 run_id: step.run_id,
1657 step_id: step.id,
1658 error,
1659 });
1660 continue;
1661 }
1662
1663 match self
1664 .store
1665 .resolve_signal_step(step.id, received_output(&stored))
1666 .await
1667 {
1668 Ok(SignalStepResolution::Resolved {
1669 run_id,
1670 run_resumed,
1671 }) => {
1672 resumed.push(SignalResumed {
1673 run_id,
1674 step_id: step.id,
1675 });
1676 if run_resumed && self.execution_mode == ExecutionMode::Local {
1677 self.spawn_local_resume(run_id);
1678 }
1679 }
1680 Ok(SignalStepResolution::NotWaiting { .. }) => {}
1682 Err(err) => {
1683 error!(
1684 run_id = %step.run_id,
1685 step_id = %step.id,
1686 error = %err,
1687 "failed to resolve a waiting signal step"
1688 );
1689 rejected.push(SignalRejected {
1690 run_id: step.run_id,
1691 step_id: step.id,
1692 error: err.to_string(),
1693 });
1694 }
1695 }
1696 }
1697
1698 info!(
1699 signal_id = %stored.id,
1700 signal = %stored.name,
1701 key = %stored.key,
1702 resumed = resumed.len(),
1703 rejected = rejected.len(),
1704 "signal received"
1705 );
1706 self.event_publisher
1707 .publish(Event::SignalReceived(SignalReceivedEvent {
1708 signal_id: stored.id,
1709 name: stored.name.clone(),
1710 key: stored.key.clone(),
1711 resumed_runs: resumed.iter().map(|r| r.run_id).collect(),
1712 at: stored.received_at,
1713 }));
1714
1715 Ok(SignalDelivery {
1716 signal_id: stored.id,
1717 duplicate: false,
1718 resumed,
1719 rejected,
1720 })
1721 }
1722
1723 pub async fn send_signal<S: Signal>(
1761 self: &Arc<Self>,
1762 signal: &S,
1763 key: &str,
1764 idempotency_id: Option<&str>,
1765 ) -> Result<SignalDelivery, EngineError> {
1766 let payload = to_value(signal)?;
1767 self.deliver_signal(NewSignal {
1768 name: S::NAME.to_string(),
1769 key: key.to_string(),
1770 payload,
1771 idempotency_id: idempotency_id.map(str::to_string),
1772 })
1773 .await
1774 }
1775
1776 pub(crate) fn spawn_local_resume(self: &Arc<Self>, run_id: Uuid) {
1782 let engine = Arc::clone(self);
1783 spawn(async move {
1784 if let Err(err) = engine
1785 .store
1786 .update_run_status(run_id, RunStatus::Running)
1787 .await
1788 {
1789 error!(run_id = %run_id, error = %err, "failed to restart a woken run");
1790 return;
1791 }
1792 if let Err(err) = engine.resume_run(run_id).await {
1793 error!(run_id = %run_id, error = %err, "failed to resume a woken run");
1794 }
1795 });
1796 }
1797
1798 pub async fn fail_or_schedule_retry(
1849 &self,
1850 run_id: Uuid,
1851 error: &str,
1852 retryable: bool,
1853 cost_usd: Option<Decimal>,
1854 duration_ms: Option<u64>,
1855 ) -> Result<RunStatus, EngineError> {
1856 let run = self
1857 .store
1858 .get_run(run_id)
1859 .await?
1860 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
1861
1862 if run.status.state == RunStatus::Paused {
1865 info!(run_id = %run_id, error = %error, "run paused, failure not recorded");
1866 return Ok(RunStatus::Paused);
1867 }
1868
1869 let has_attempts_left = run.retry_count < run.max_retries;
1870 let update = if retryable && has_attempts_left {
1871 let backoff = backoff_for_retry(run.retry_count);
1872 let scheduled_at = Utc::now() + TimeDelta::milliseconds(backoff.as_millis() as i64);
1873
1874 info!(
1875 run_id = %run_id,
1876 workflow = %run.workflow_name,
1877 attempt = run.retry_count + 1,
1878 max_retries = run.max_retries,
1879 backoff_secs = backoff.as_secs(),
1880 scheduled_at = %scheduled_at,
1881 "run failed, scheduling retry"
1882 );
1883
1884 RunUpdate {
1885 status: Some(RunStatus::Retrying),
1886 error: Some(error.to_string()),
1887 increment_retry: true,
1888 cost_usd,
1889 duration_ms,
1890 scheduled_at: Some(scheduled_at),
1891 ..RunUpdate::default()
1892 }
1893 } else {
1894 RunUpdate {
1895 status: Some(RunStatus::Failed),
1896 error: Some(error.to_string()),
1897 cost_usd,
1898 duration_ms,
1899 completed_at: Some(Utc::now()),
1900 ..RunUpdate::default()
1901 }
1902 };
1903
1904 let status = update.status.unwrap_or(RunStatus::Failed);
1905 self.store.update_run(run_id, update).await?;
1906 self.fail_orphaned_steps(run_id, error).await?;
1907 self.cancel_descendants_of_stopped_run(run_id, error).await;
1910
1911 Ok(status)
1912 }
1913
1914 pub async fn interrupt_running_steps(&self, run_id: Uuid) -> Result<(), EngineError> {
1944 interrupt_running_steps(self.store.as_ref(), run_id).await
1945 }
1946
1947 pub async fn fail_orphaned_steps(
1961 &self,
1962 run_id: Uuid,
1963 error_message: &str,
1964 ) -> Result<(), EngineError> {
1965 let steps = self.store.list_steps(run_id).await?;
1966 let now = Utc::now();
1967
1968 for step in steps {
1969 if step.status.state.is_terminal() {
1970 continue;
1971 }
1972
1973 let (target_status, error) = match step.status.state {
1974 StepStatus::Running | StepStatus::AwaitingApproval => {
1975 let err = if step.error.is_some() {
1976 None
1977 } else {
1978 Some(error_message.to_string())
1979 };
1980 (StepStatus::Failed, err)
1981 }
1982 StepStatus::Pending => (StepStatus::Skipped, None),
1983 _ => continue,
1984 };
1985
1986 if let Err(e) = self
1987 .store
1988 .update_step(
1989 step.id,
1990 StepUpdate {
1991 status: Some(target_status),
1992 error,
1993 completed_at: Some(now),
1994 ..StepUpdate::default()
1995 },
1996 )
1997 .await
1998 {
1999 warn!(
2000 run_id = %run_id,
2001 step_id = %step.id,
2002 step_name = %step.name,
2003 error = %e,
2004 "failed to cleanup orphaned step"
2005 );
2006 } else {
2007 info!(
2008 run_id = %run_id,
2009 step_id = %step.id,
2010 step_name = %step.name,
2011 from = %step.status.state,
2012 to = %target_status,
2013 "cleaned up orphaned step"
2014 );
2015 }
2016 }
2017
2018 Ok(())
2019 }
2020
2021 async fn release_then_execute(
2026 &self,
2027 run_id: Uuid,
2028 handler: &dyn WorkflowHandler,
2029 ctx: &mut WorkflowContext,
2030 ) -> Result<(), EngineError> {
2031 match self.provider.release_run(&run_id.to_string()).await {
2032 Ok(()) => handler.execute(ctx).await,
2033 Err(e) => Err(EngineError::Operation(OperationError::Agent(e))),
2034 }
2035 }
2036
2037 async fn finalize_run(
2043 &self,
2044 run_id: Uuid,
2045 workflow_name: &str,
2046 result: Result<(), EngineError>,
2047 ctx: &WorkflowContext,
2048 run_start: Instant,
2049 run_labels: HashMap<String, String>,
2050 ) -> Result<WorkflowResult, EngineError> {
2051 let total_duration = ctx.carried_duration_ms() + run_start.elapsed().as_millis() as u64;
2054 let completed_at = Utc::now();
2055
2056 if let Some(run) = self.store.get_run(run_id).await?
2060 && run.status.state == RunStatus::Paused
2061 {
2062 let run = self
2063 .store
2064 .update_run_returning(
2065 run_id,
2066 RunUpdate {
2067 cost_usd: Some(ctx.total_cost_usd()),
2068 duration_ms: Some(total_duration),
2069 ..RunUpdate::default()
2070 },
2071 )
2072 .await?;
2073 info!(
2074 run_id = %run_id,
2075 outcome = ?result.err().map(|err| err.to_string()),
2076 "run paused, execution stopped"
2077 );
2078 return Ok(WorkflowResult {
2079 run,
2080 steps: ctx.step_results().to_vec(),
2081 });
2082 }
2083
2084 let final_status;
2085 let final_run;
2086
2087 match result {
2088 Ok(()) => {
2089 final_status = if ctx.has_allowed_failure() {
2090 RunStatus::Warning
2091 } else {
2092 RunStatus::Completed
2093 };
2094 final_run = self
2095 .store
2096 .update_run_returning(
2097 run_id,
2098 RunUpdate {
2099 status: Some(final_status),
2100 cost_usd: Some(ctx.total_cost_usd()),
2101 duration_ms: Some(total_duration),
2102 completed_at: Some(completed_at),
2103 output: ctx.output().cloned(),
2104 ..RunUpdate::default()
2105 },
2106 )
2107 .await?;
2108
2109 info!(
2110 run_id = %run_id,
2111 status = %final_status,
2112 cost_usd = %ctx.total_cost_usd(),
2113 duration_ms = total_duration,
2114 "run completed"
2115 );
2116 }
2117 Err(EngineError::ApprovalRequired {
2118 run_id: approval_run_id,
2119 step_id,
2120 ref message,
2121 }) => {
2122 final_status = RunStatus::AwaitingApproval;
2123 final_run = self
2124 .store
2125 .update_run_returning(
2126 run_id,
2127 RunUpdate {
2128 status: Some(RunStatus::AwaitingApproval),
2129 cost_usd: Some(ctx.total_cost_usd()),
2130 duration_ms: Some(total_duration),
2131 ..RunUpdate::default()
2132 },
2133 )
2134 .await?;
2135
2136 info!(
2137 run_id = %approval_run_id,
2138 step_id = %step_id,
2139 message = %message,
2140 "run awaiting approval"
2141 );
2142
2143 self.publish_approval_requested(approval_run_id, step_id, message)
2144 .await?;
2145 }
2146 Err(EngineError::ChildSuspended {
2147 run_id: child_run_id,
2148 ref cause,
2149 }) => {
2150 final_status = cause.suspension_status();
2151 final_run = self
2155 .store
2156 .update_run_returning(
2157 run_id,
2158 RunUpdate {
2159 status: Some(final_status),
2160 cost_usd: Some(ctx.total_cost_usd()),
2161 duration_ms: Some(total_duration),
2162 ..RunUpdate::default()
2163 },
2164 )
2165 .await?;
2166
2167 let leaf = cause.suspension_leaf();
2168 info!(
2169 run_id = %run_id,
2170 child_run_id = %child_run_id,
2171 status = %final_status,
2172 cause = %leaf,
2173 "run suspended with its child run"
2174 );
2175
2176 match leaf {
2177 EngineError::ApprovalRequired {
2178 run_id: approval_run_id,
2179 step_id,
2180 message,
2181 } => {
2182 self.publish_approval_requested(*approval_run_id, *step_id, message)
2183 .await?;
2184 }
2185 EngineError::SignalWaiting {
2186 run_id: wait_run_id,
2187 step_id,
2188 step_name,
2189 name,
2190 key,
2191 deadline_at,
2192 } => {
2193 self.event_publisher
2194 .publish(Event::SignalAwaited(SignalAwaitedEvent {
2195 run_id: *wait_run_id,
2196 step_id: *step_id,
2197 step_name: step_name.clone(),
2198 name: name.clone(),
2199 key: key.clone(),
2200 deadline_at: *deadline_at,
2201 at: Utc::now(),
2202 }));
2203 }
2204 _ => {}
2207 }
2208 }
2209 Err(EngineError::HumanInputRequired {
2210 run_id: input_run_id,
2211 step_id,
2212 ref message,
2213 }) => {
2214 final_status = RunStatus::AwaitingApproval;
2215 final_run = self
2216 .store
2217 .update_run_returning(
2218 run_id,
2219 RunUpdate {
2220 status: Some(RunStatus::AwaitingApproval),
2221 cost_usd: Some(ctx.total_cost_usd()),
2222 duration_ms: Some(total_duration),
2223 ..RunUpdate::default()
2224 },
2225 )
2226 .await?;
2227
2228 info!(
2230 run_id = %input_run_id,
2231 step_id = %step_id,
2232 message = %message,
2233 "run awaiting human input"
2234 );
2235 }
2236 Err(EngineError::DelaySleeping {
2237 run_id: delay_run_id,
2238 step_id,
2239 wake_at,
2240 }) => {
2241 final_status = RunStatus::Sleeping;
2242 final_run = self
2243 .store
2244 .update_run_returning(
2245 run_id,
2246 RunUpdate {
2247 status: Some(RunStatus::Sleeping),
2248 cost_usd: Some(ctx.total_cost_usd()),
2249 duration_ms: Some(total_duration),
2250 scheduled_at: Some(wake_at),
2251 ..RunUpdate::default()
2252 },
2253 )
2254 .await?;
2255
2256 info!(
2257 run_id = %delay_run_id,
2258 step_id = %step_id,
2259 wake_at = %wake_at,
2260 "run sleeping until delay elapses"
2261 );
2262 }
2263 Err(EngineError::CapacitySleeping {
2264 run_id: capacity_run_id,
2265 step_id,
2266 ref kind,
2267 wake_at,
2268 }) => {
2269 final_status = RunStatus::Sleeping;
2270 final_run = self
2271 .store
2272 .update_run_returning(
2273 run_id,
2274 RunUpdate {
2275 status: Some(RunStatus::Sleeping),
2276 cost_usd: Some(ctx.total_cost_usd()),
2277 duration_ms: Some(total_duration),
2278 scheduled_at: Some(wake_at),
2279 capacity_wait_kind: Some(ProviderKind::new(kind.as_str())),
2280 ..RunUpdate::default()
2281 },
2282 )
2283 .await?;
2284
2285 info!(
2286 run_id = %capacity_run_id,
2287 step_id = %step_id,
2288 kind = %kind,
2289 wake_at = %wake_at,
2290 "run sleeping until provider capacity returns"
2291 );
2292 }
2293 Err(EngineError::SignalWaiting {
2294 run_id: wait_run_id,
2295 step_id,
2296 ref step_name,
2297 ref name,
2298 ref key,
2299 deadline_at,
2300 }) => {
2301 final_status = RunStatus::Sleeping;
2302 let waiting = self
2306 .store
2307 .suspend_run_on_signal(run_id, step_id, deadline_at)
2308 .await?;
2309 final_run = self
2310 .store
2311 .update_run_returning(
2312 run_id,
2313 RunUpdate {
2314 cost_usd: Some(ctx.total_cost_usd()),
2315 duration_ms: Some(total_duration),
2316 ..RunUpdate::default()
2317 },
2318 )
2319 .await?;
2320
2321 if waiting {
2322 self.event_publisher
2323 .publish(Event::SignalAwaited(SignalAwaitedEvent {
2324 run_id: wait_run_id,
2325 step_id,
2326 step_name: step_name.clone(),
2327 name: name.clone(),
2328 key: key.clone(),
2329 deadline_at,
2330 at: Utc::now(),
2331 }));
2332 }
2333
2334 info!(
2335 run_id = %wait_run_id,
2336 step_id = %step_id,
2337 signal = %name,
2338 key = %key,
2339 deadline_at = %deadline_at,
2340 waiting,
2341 "run sleeping until a signal arrives"
2342 );
2343 }
2344 Err(err) => {
2345 let guardrail_stop = matches!(
2349 err,
2350 EngineError::RunBudgetExceeded { .. } | EngineError::WorkflowGuardRejected(_)
2351 );
2352
2353 final_status = if guardrail_stop {
2354 if let Err(store_err) = self
2355 .store
2356 .update_run(
2357 run_id,
2358 RunUpdate {
2359 status: Some(RunStatus::Cancelled),
2360 error: Some(err.to_string()),
2361 cost_usd: Some(ctx.total_cost_usd()),
2362 duration_ms: Some(total_duration),
2363 completed_at: Some(completed_at),
2364 output: ctx.output().cloned(),
2365 ..RunUpdate::default()
2366 },
2367 )
2368 .await
2369 {
2370 error!(run_id = %run_id, store_error = %store_err, "failed to persist run cancellation");
2371 }
2372 if let Err(cleanup_err) = self
2373 .fail_orphaned_steps(run_id, "run stopped: guardrail limit reached")
2374 .await
2375 {
2376 error!(run_id = %run_id, store_error = %cleanup_err, "failed to cleanup orphaned steps");
2377 }
2378 RunStatus::Cancelled
2379 } else {
2380 if let Some(output) = ctx.output()
2383 && let Err(store_err) = self
2384 .store
2385 .update_run(
2386 run_id,
2387 RunUpdate {
2388 output: Some(output.clone()),
2389 ..RunUpdate::default()
2390 },
2391 )
2392 .await
2393 {
2394 error!(run_id = %run_id, store_error = %store_err, "failed to persist run output");
2395 }
2396 self.fail_or_schedule_retry(
2397 run_id,
2398 &err.to_string(),
2399 is_run_retryable(&err),
2400 Some(ctx.total_cost_usd()),
2401 Some(total_duration),
2402 )
2403 .await
2404 .unwrap_or_else(|store_err| {
2405 error!(run_id = %run_id, store_error = %store_err, "failed to persist run failure");
2406 RunStatus::Failed
2407 })
2408 };
2409
2410 if matches!(err, EngineError::RunBudgetExceeded { .. }) {
2411 self.on_run_budget_exceeded(workflow_name, run_id, &err);
2412 }
2413
2414 error!(run_id = %run_id, status = %final_status, error = %err, "run stopped");
2415
2416 self.publish_run_status_changed(
2417 workflow_name,
2418 run_id,
2419 final_status,
2420 Some(err.to_string()),
2421 ctx,
2422 total_duration,
2423 run_labels,
2424 );
2425
2426 #[cfg(feature = "prometheus")]
2427 self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
2428
2429 return Err(err);
2430 }
2431 }
2432
2433 self.publish_run_status_changed(
2434 workflow_name,
2435 run_id,
2436 final_status,
2437 None,
2438 ctx,
2439 total_duration,
2440 run_labels,
2441 );
2442
2443 #[cfg(feature = "prometheus")]
2444 self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
2445
2446 Ok(WorkflowResult {
2447 run: final_run,
2448 steps: ctx.step_results().to_vec(),
2449 })
2450 }
2451
2452 async fn publish_approval_requested(
2455 &self,
2456 run_id: Uuid,
2457 step_id: Uuid,
2458 message: &str,
2459 ) -> Result<(), EngineError> {
2460 let requirement = self
2461 .store
2462 .get_step(step_id)
2463 .await?
2464 .and_then(|s| s.approval_requirement);
2465 self.event_publisher
2466 .publish(Event::ApprovalRequested(ApprovalRequestedEvent {
2467 run_id,
2468 step_id,
2469 message: message.to_string(),
2470 requirement,
2471 at: Utc::now(),
2472 }));
2473 Ok(())
2474 }
2475
2476 pub async fn fail_ancestors(&self, run_id: Uuid, reason: &str) -> Result<(), EngineError> {
2509 let mut current = self
2510 .store
2511 .get_run(run_id)
2512 .await?
2513 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
2514 let mut visited = HashSet::from([run_id]);
2516
2517 while let Some(parent_id) = chain_parent(¤t) {
2518 if !visited.insert(parent_id) {
2519 break;
2520 }
2521 let status = self
2522 .fail_or_schedule_retry(parent_id, reason, false, None, None)
2523 .await?;
2524 if status == RunStatus::Paused {
2527 self.requeue_paused_root(¤t).await?;
2528 break;
2529 }
2530 info!(
2531 run_id = %run_id,
2532 ancestor_run_id = %parent_id,
2533 status = %status,
2534 "ancestor run failed with its child"
2535 );
2536 current = self
2537 .store
2538 .get_run(parent_id)
2539 .await?
2540 .ok_or(EngineError::Store(StoreError::RunNotFound(parent_id)))?;
2541 }
2542
2543 Ok(())
2544 }
2545
2546 #[cfg(feature = "prometheus")]
2548 fn emit_run_metrics(
2549 &self,
2550 workflow_name: &str,
2551 status: RunStatus,
2552 duration_ms: u64,
2553 ctx: &WorkflowContext,
2554 ) {
2555 let status_str = status.to_string();
2556 let wf = workflow_name.to_string();
2557
2558 counter!(RUNS_TOTAL, "workflow" => wf.clone(), "status" => status_str.clone()).increment(1);
2559 histogram!(RUN_DURATION_SECONDS, "workflow" => wf.clone(), "status" => status_str)
2560 .record(duration_ms as f64 / 1000.0);
2561 histogram!(RUN_COST_USD, "workflow" => wf.clone()).record(
2562 ctx.total_cost_usd()
2563 .to_string()
2564 .parse::<f64>()
2565 .unwrap_or(0.0),
2566 );
2567 gauge!(RUNS_ACTIVE, "workflow" => wf).decrement(1.0);
2568 }
2569
2570 fn on_run_budget_exceeded(&self, workflow_name: &str, run_id: Uuid, err: &EngineError) {
2576 let EngineError::RunBudgetExceeded {
2577 limit_usd,
2578 spent_usd,
2579 step_budget_usd,
2580 ..
2581 } = err
2582 else {
2583 return;
2584 };
2585
2586 #[cfg(feature = "prometheus")]
2587 counter!(
2588 RUN_BUDGET_EXCEEDED_TOTAL,
2589 "workflow" => workflow_name.to_string(),
2590 "scope" => "run",
2591 )
2592 .increment(1);
2593
2594 self.event_publisher
2595 .publish(Event::RunBudgetExceeded(RunBudgetExceededEvent {
2596 run_id,
2597 workflow_name: workflow_name.to_string(),
2598 limit_usd: *limit_usd,
2599 spent_usd: *spent_usd,
2600 step_budget_usd: *step_budget_usd,
2601 at: Utc::now(),
2602 }));
2603 }
2604
2605 #[allow(clippy::too_many_arguments)]
2610 fn publish_run_status_changed(
2611 &self,
2612 workflow_name: &str,
2613 run_id: Uuid,
2614 to: RunStatus,
2615 error: Option<String>,
2616 ctx: &WorkflowContext,
2617 duration_ms: u64,
2618 labels: HashMap<String, String>,
2619 ) {
2620 let now = Utc::now();
2621 let cost_usd = ctx.total_cost_usd();
2622 let wf = workflow_name.to_string();
2623
2624 self.event_publisher
2625 .publish(Event::RunStatusChanged(RunStatusChangedEvent {
2626 run_id,
2627 workflow_name: wf.clone(),
2628 from: RunStatus::Running,
2629 to,
2630 error: error.clone(),
2631 cost_usd,
2632 duration_ms,
2633 labels: labels.clone(),
2634 at: now,
2635 }));
2636
2637 if to == RunStatus::Failed {
2638 self.event_publisher
2639 .publish(Event::RunFailed(RunFailedEvent {
2640 run_id,
2641 workflow_name: wf,
2642 error,
2643 cost_usd,
2644 duration_ms,
2645 labels,
2646 at: now,
2647 }));
2648 }
2649 }
2650}
2651
2652impl fmt::Debug for Engine {
2653 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
2654 f.debug_struct("Engine")
2655 .field("handlers", &self.handlers.keys().collect::<Vec<_>>())
2656 .finish_non_exhaustive()
2657 }
2658}
2659
2660#[cfg(test)]
2661mod tests {
2662 use super::*;
2663 use crate::config::ShellConfig;
2664 use crate::handler::{HandlerFuture, WorkflowHandler};
2665 use ironflow_core::providers::claude::ClaudeCodeProvider;
2666 use ironflow_core::providers::record_replay::RecordReplayProvider;
2667 use ironflow_store::memory::InMemoryStore;
2668 use ironflow_store::models::{MAX_PRIORITY, MIN_PRIORITY, StepStatus};
2669 use serde_json::json;
2670
2671 struct EchoWorkflow;
2673
2674 impl WorkflowHandler for EchoWorkflow {
2675 fn name(&self) -> &str {
2676 "echo-workflow"
2677 }
2678
2679 fn describe(&self) -> WorkflowInfo {
2680 WorkflowInfo {
2681 description: "A simple workflow that echoes hello".to_string(),
2682 source_code: None,
2683 sub_workflows: Vec::new(),
2684 category: None,
2685 version: self.version().map(str::to_string),
2686 compatible_versions: Vec::new(),
2687 input_schema: None,
2688 default_labels: HashMap::new(),
2689 schedule: self.schedule().cloned(),
2690 default_max_cost_usd: self.default_max_cost_usd(),
2691 priority: self.priority(),
2692 }
2693 }
2694
2695 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2696 Box::pin(async move {
2697 ctx.shell("greet", ShellConfig::new("echo hello")).await?;
2698 Ok(())
2699 })
2700 }
2701 }
2702
2703 struct FailingWorkflow;
2705
2706 impl WorkflowHandler for FailingWorkflow {
2707 fn name(&self) -> &str {
2708 "failing-workflow"
2709 }
2710
2711 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2712 Box::pin(async move {
2713 ctx.shell("fail", ShellConfig::new("exit 1")).await?;
2714 Ok(())
2715 })
2716 }
2717 }
2718
2719 fn create_test_engine() -> Engine {
2720 let store = Arc::new(InMemoryStore::new());
2721 let inner = ClaudeCodeProvider::new();
2722 let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
2723 inner,
2724 "/tmp/ironflow-fixtures",
2725 ));
2726 Engine::new(store, provider)
2727 }
2728
2729 #[test]
2730 fn engine_new_creates_instance() {
2731 let engine = create_test_engine();
2732 assert_eq!(engine.handler_names().len(), 0);
2733 }
2734
2735 #[test]
2736 fn execution_mode_defaults_to_local() {
2737 let engine = create_test_engine();
2738 assert_eq!(engine.execution_mode(), ExecutionMode::Local);
2739 }
2740
2741 #[test]
2742 fn with_execution_mode_overrides_the_default() {
2743 let engine = create_test_engine().with_execution_mode(ExecutionMode::Workers);
2744 assert_eq!(engine.execution_mode(), ExecutionMode::Workers);
2745 }
2746
2747 #[test]
2748 fn engine_register_handler() {
2749 let mut engine = create_test_engine();
2750 let result = engine.register(EchoWorkflow);
2751 assert!(result.is_ok());
2752 assert_eq!(engine.handler_names().len(), 1);
2753 assert!(engine.handler_names().contains(&"echo-workflow"));
2754 }
2755
2756 #[test]
2757 fn engine_register_duplicate_returns_error() {
2758 let mut engine = create_test_engine();
2759 engine.register(EchoWorkflow).unwrap();
2760 let result = engine.register(EchoWorkflow);
2761 assert!(result.is_err());
2762 }
2763
2764 #[test]
2765 fn engine_get_handler_found() {
2766 let mut engine = create_test_engine();
2767 engine.register(EchoWorkflow).unwrap();
2768 let handler = engine.get_handler("echo-workflow");
2769 assert!(handler.is_some());
2770 }
2771
2772 #[test]
2773 fn engine_get_handler_not_found() {
2774 let engine = create_test_engine();
2775 let handler = engine.get_handler("nonexistent");
2776 assert!(handler.is_none());
2777 }
2778
2779 #[test]
2780 fn engine_handler_names_lists_all() {
2781 let mut engine = create_test_engine();
2782 engine.register(EchoWorkflow).unwrap();
2783 engine.register(FailingWorkflow).unwrap();
2784 let names = engine.handler_names();
2785 assert_eq!(names.len(), 2);
2786 assert!(names.contains(&"echo-workflow"));
2787 assert!(names.contains(&"failing-workflow"));
2788 }
2789
2790 #[test]
2791 fn engine_handler_info_returns_description() {
2792 let mut engine = create_test_engine();
2793 engine.register(EchoWorkflow).unwrap();
2794 let info = engine.handler_info("echo-workflow");
2795 assert!(info.is_some());
2796 let info = info.unwrap();
2797 assert_eq!(info.description, "A simple workflow that echoes hello");
2798 }
2799
2800 struct CategorizedWorkflow;
2801
2802 impl WorkflowHandler for CategorizedWorkflow {
2803 fn name(&self) -> &str {
2804 "categorized"
2805 }
2806 fn category(&self) -> Option<&str> {
2807 Some("data/etl")
2808 }
2809 fn execute<'a>(
2810 &'a self,
2811 _ctx: &'a mut WorkflowContext,
2812 ) -> crate::handler::HandlerFuture<'a> {
2813 Box::pin(async move { Ok(()) })
2814 }
2815 }
2816
2817 #[test]
2818 fn engine_default_describe_propagates_category() {
2819 let mut engine = create_test_engine();
2820 engine.register(CategorizedWorkflow).unwrap();
2821 let info = engine.handler_info("categorized").unwrap();
2822 assert_eq!(info.category.as_deref(), Some("data/etl"));
2823 }
2824
2825 #[test]
2826 fn engine_default_describe_without_category() {
2827 let mut engine = create_test_engine();
2828 engine.register(EchoWorkflow).unwrap();
2829 let info = engine.handler_info("echo-workflow").unwrap();
2830 assert!(info.category.is_none());
2831 }
2832
2833 struct ScheduledWorkflow {
2838 schedule: CronSchedule,
2839 }
2840
2841 impl ScheduledWorkflow {
2842 fn new() -> Self {
2843 Self {
2844 schedule: CronSchedule::new("0 0 * * * *").unwrap(),
2845 }
2846 }
2847 }
2848
2849 impl WorkflowHandler for ScheduledWorkflow {
2850 fn name(&self) -> &str {
2851 "scheduled"
2852 }
2853 fn schedule(&self) -> Option<&CronSchedule> {
2854 Some(&self.schedule)
2855 }
2856 fn execute<'a>(
2857 &'a self,
2858 _ctx: &'a mut WorkflowContext,
2859 ) -> crate::handler::HandlerFuture<'a> {
2860 Box::pin(async move { Ok(()) })
2861 }
2862 }
2863
2864 #[test]
2865 fn engine_default_describe_propagates_schedule() {
2866 let mut engine = create_test_engine();
2867 engine.register(ScheduledWorkflow::new()).unwrap();
2868 let info = engine.handler_info("scheduled").unwrap();
2869 assert_eq!(
2870 info.schedule.as_ref().map(|s| s.as_str()),
2871 Some("0 0 * * * *")
2872 );
2873 }
2874
2875 #[test]
2876 fn engine_default_describe_without_schedule() {
2877 let mut engine = create_test_engine();
2878 engine.register(EchoWorkflow).unwrap();
2879 let info = engine.handler_info("echo-workflow").unwrap();
2880 assert!(info.schedule.is_none());
2881 }
2882
2883 #[test]
2884 fn scheduled_handlers_returns_only_scheduled() {
2885 let mut engine = create_test_engine();
2886 engine.register(EchoWorkflow).unwrap();
2887 engine.register(ScheduledWorkflow::new()).unwrap();
2888 engine.register(FailingWorkflow).unwrap();
2889
2890 let scheduled = engine.scheduled_handlers();
2891 assert_eq!(scheduled.len(), 1);
2892 assert_eq!(scheduled[0].0, "scheduled");
2893 assert_eq!(scheduled[0].1.as_str(), "0 0 * * * *");
2894 }
2895
2896 #[test]
2897 fn scheduled_handlers_empty_when_none_scheduled() {
2898 let mut engine = create_test_engine();
2899 engine.register(EchoWorkflow).unwrap();
2900 engine.register(FailingWorkflow).unwrap();
2901
2902 let scheduled = engine.scheduled_handlers();
2903 assert!(scheduled.is_empty());
2904 }
2905
2906 struct BadCategoryWorkflow(&'static str);
2907
2908 impl WorkflowHandler for BadCategoryWorkflow {
2909 fn name(&self) -> &str {
2910 "bad-category"
2911 }
2912 fn category(&self) -> Option<&str> {
2913 Some(self.0)
2914 }
2915 fn execute<'a>(
2916 &'a self,
2917 _ctx: &'a mut WorkflowContext,
2918 ) -> crate::handler::HandlerFuture<'a> {
2919 Box::pin(async move { Ok(()) })
2920 }
2921 }
2922
2923 #[test]
2924 fn engine_register_rejects_empty_category() {
2925 let mut engine = create_test_engine();
2926 let err = engine.register(BadCategoryWorkflow("")).unwrap_err();
2927 match err {
2928 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty category")),
2929 other => panic!("expected InvalidWorkflow, got {other:?}"),
2930 }
2931 }
2932
2933 #[test]
2934 fn engine_register_rejects_leading_slash_category() {
2935 let mut engine = create_test_engine();
2936 let err = engine
2937 .register(BadCategoryWorkflow("/data/etl"))
2938 .unwrap_err();
2939 match err {
2940 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("leading '/'")),
2941 other => panic!("expected InvalidWorkflow, got {other:?}"),
2942 }
2943 }
2944
2945 #[test]
2946 fn engine_register_rejects_trailing_slash_category() {
2947 let mut engine = create_test_engine();
2948 let err = engine
2949 .register(BadCategoryWorkflow("data/etl/"))
2950 .unwrap_err();
2951 match err {
2952 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("trailing '/'")),
2953 other => panic!("expected InvalidWorkflow, got {other:?}"),
2954 }
2955 }
2956
2957 #[test]
2958 fn engine_register_rejects_double_slash_category() {
2959 let mut engine = create_test_engine();
2960 let err = engine
2961 .register(BadCategoryWorkflow("data//etl"))
2962 .unwrap_err();
2963 match err {
2964 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty segment")),
2965 other => panic!("expected InvalidWorkflow, got {other:?}"),
2966 }
2967 }
2968
2969 #[test]
2970 fn engine_register_rejects_whitespace_only_segment_category() {
2971 let mut engine = create_test_engine();
2972 let err = engine
2973 .register(BadCategoryWorkflow("data/ /etl"))
2974 .unwrap_err();
2975 match err {
2976 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("whitespace-only segment")),
2977 other => panic!("expected InvalidWorkflow, got {other:?}"),
2978 }
2979 }
2980
2981 #[test]
2982 fn engine_register_accepts_valid_nested_category() {
2983 let mut engine = create_test_engine();
2984 assert!(engine.register(CategorizedWorkflow).is_ok());
2985 }
2986
2987 #[tokio::test]
2988 async fn engine_unknown_workflow_returns_error() {
2989 let engine = create_test_engine();
2990 let result = engine
2991 .run_handler("unknown", TriggerKind::Manual, json!({}))
2992 .await;
2993 assert!(result.is_err());
2994 match result {
2995 Err(EngineError::InvalidWorkflow(msg)) => {
2996 assert!(msg.contains("no handler registered"));
2997 }
2998 _ => panic!("expected InvalidWorkflow error"),
2999 }
3000 }
3001
3002 #[tokio::test]
3003 async fn engine_enqueue_handler_creates_pending_run() {
3004 let mut engine = create_test_engine();
3005 engine.register(EchoWorkflow).unwrap();
3006
3007 let run = engine
3008 .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
3009 .await
3010 .unwrap();
3011 assert_eq!(run.status.state, RunStatus::Pending);
3012 assert_eq!(run.workflow_name, "echo-workflow");
3013 }
3014
3015 #[tokio::test]
3016 async fn enqueue_handler_leaves_the_run_unattributed() {
3017 let mut engine = create_test_engine();
3018 engine.register(EchoWorkflow).unwrap();
3019
3020 let run = engine
3021 .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
3022 .await
3023 .unwrap();
3024
3025 assert!(run.created_by.is_none());
3026 }
3027
3028 #[tokio::test]
3029 async fn enqueue_handler_with_options_records_the_author() {
3030 let mut engine = create_test_engine();
3031 engine.register(EchoWorkflow).unwrap();
3032 let actor = RunActor::User {
3033 user_id: Uuid::now_v7(),
3034 };
3035
3036 let run = engine
3037 .enqueue_handler_with_options(
3038 "echo-workflow",
3039 TriggerKind::Api,
3040 json!({}),
3041 EnqueueOptions {
3042 created_by: Some(actor.clone()),
3043 ..Default::default()
3044 },
3045 )
3046 .await
3047 .unwrap()
3048 .into_run();
3049
3050 assert_eq!(run.created_by, Some(actor));
3051 }
3052
3053 #[tokio::test]
3054 async fn enqueue_handler_with_options_accepts_no_author() {
3055 let mut engine = create_test_engine();
3056 engine.register(EchoWorkflow).unwrap();
3057
3058 let run = engine
3059 .enqueue_handler_with_options(
3060 "echo-workflow",
3061 TriggerKind::Cron {
3062 schedule: "0 * * * * *".to_string(),
3063 schedule_id: None,
3064 scheduled_for: None,
3065 },
3066 json!({}),
3067 EnqueueOptions::default(),
3068 )
3069 .await
3070 .unwrap()
3071 .into_run();
3072
3073 assert!(run.created_by.is_none());
3074 }
3075
3076 #[tokio::test]
3077 async fn enqueue_handler_with_options_stores_concurrency_limits() {
3078 let mut engine = create_test_engine();
3079 engine.register(EchoWorkflow).unwrap();
3080 let limits = vec![
3081 ConcurrencyLimit::new("repo:acme", 2),
3082 ConcurrencyLimit::new("tenant:42", 5),
3083 ];
3084
3085 let run = engine
3086 .enqueue_handler_with_options(
3087 "echo-workflow",
3088 TriggerKind::Api,
3089 json!({}),
3090 EnqueueOptions {
3091 concurrency_limits: limits.clone(),
3092 ..Default::default()
3093 },
3094 )
3095 .await
3096 .unwrap()
3097 .into_run();
3098
3099 assert_eq!(run.concurrency_limits, limits);
3100 }
3101
3102 #[tokio::test]
3103 async fn enqueue_rejects_invalid_concurrency_limits() {
3104 let mut engine = create_test_engine();
3105 engine.register(EchoWorkflow).unwrap();
3106
3107 let invalid = [
3108 vec![ConcurrencyLimit::new("repo:acme", 0)],
3109 vec![ConcurrencyLimit::new("", 1)],
3110 vec![
3111 ConcurrencyLimit::new("repo:acme", 1),
3112 ConcurrencyLimit::new("repo:acme", 2),
3113 ],
3114 ];
3115 for concurrency_limits in invalid {
3116 let err = engine
3117 .enqueue_handler_with_options(
3118 "echo-workflow",
3119 TriggerKind::Api,
3120 json!({}),
3121 EnqueueOptions {
3122 concurrency_limits,
3123 ..Default::default()
3124 },
3125 )
3126 .await
3127 .unwrap_err();
3128 assert!(
3129 matches!(err, EngineError::InvalidConcurrencyLimit(_)),
3130 "{err:?}"
3131 );
3132 }
3133
3134 let err = engine
3136 .enqueue_handler_with_options(
3137 "not-registered",
3138 TriggerKind::Api,
3139 json!({}),
3140 EnqueueOptions {
3141 concurrency_limits: vec![ConcurrencyLimit::new("repo:acme", 0)],
3142 ..Default::default()
3143 },
3144 )
3145 .await
3146 .unwrap_err();
3147 assert!(
3148 matches!(err, EngineError::InvalidConcurrencyLimit(_)),
3149 "{err:?}"
3150 );
3151
3152 let page = engine
3153 .store()
3154 .list_runs(RunFilter::default(), 1, 10)
3155 .await
3156 .unwrap();
3157 assert_eq!(page.total, 0, "no run may be created");
3158 }
3159
3160 struct UrgentWorkflow;
3161
3162 impl WorkflowHandler for UrgentWorkflow {
3163 fn name(&self) -> &str {
3164 "urgent-workflow"
3165 }
3166
3167 fn priority(&self) -> i16 {
3168 60
3169 }
3170
3171 fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3172 Box::pin(async { Ok(()) })
3173 }
3174 }
3175
3176 #[tokio::test]
3177 async fn enqueue_priority_defaults_to_the_handler_priority() {
3178 let mut engine = create_test_engine();
3179 engine.register(EchoWorkflow).unwrap();
3180 engine.register(UrgentWorkflow).unwrap();
3181
3182 let echo = engine
3183 .enqueue_handler_with_options(
3184 "echo-workflow",
3185 TriggerKind::Api,
3186 json!({}),
3187 EnqueueOptions::default(),
3188 )
3189 .await
3190 .unwrap()
3191 .into_run();
3192 assert_eq!(echo.priority, 0);
3193
3194 let urgent = engine
3195 .enqueue_handler_with_options(
3196 "urgent-workflow",
3197 TriggerKind::Api,
3198 json!({}),
3199 EnqueueOptions::default(),
3200 )
3201 .await
3202 .unwrap()
3203 .into_run();
3204 assert_eq!(urgent.priority, 60);
3205 }
3206
3207 #[tokio::test]
3208 async fn enqueue_priority_explicit_value_overrides_the_handler() {
3209 let mut engine = create_test_engine();
3210 engine.register(UrgentWorkflow).unwrap();
3211
3212 let run = engine
3213 .enqueue_handler_with_options(
3214 "urgent-workflow",
3215 TriggerKind::Api,
3216 json!({}),
3217 EnqueueOptions {
3218 priority: Some(-20),
3219 ..Default::default()
3220 },
3221 )
3222 .await
3223 .unwrap()
3224 .into_run();
3225 assert_eq!(run.priority, -20);
3226
3227 let stored = engine.store().get_run(run.id).await.unwrap().unwrap();
3228 assert_eq!(stored.priority, -20);
3229 }
3230
3231 #[tokio::test]
3232 async fn enqueue_priority_out_of_range_is_rejected() {
3233 let mut engine = create_test_engine();
3234 engine.register(EchoWorkflow).unwrap();
3235
3236 for priority in [MAX_PRIORITY + 1, MIN_PRIORITY - 1] {
3237 let err = engine
3238 .enqueue_handler_with_options(
3239 "echo-workflow",
3240 TriggerKind::Api,
3241 json!({}),
3242 EnqueueOptions {
3243 priority: Some(priority),
3244 ..Default::default()
3245 },
3246 )
3247 .await
3248 .unwrap_err();
3249 assert!(matches!(err, EngineError::InvalidPriority(_)), "{err:?}");
3250 }
3251
3252 let err = engine
3254 .enqueue_handler_with_options(
3255 "not-registered",
3256 TriggerKind::Api,
3257 json!({}),
3258 EnqueueOptions {
3259 priority: Some(MAX_PRIORITY + 1),
3260 ..Default::default()
3261 },
3262 )
3263 .await
3264 .unwrap_err();
3265 assert!(matches!(err, EngineError::InvalidPriority(_)), "{err:?}");
3266
3267 let page = engine
3268 .store()
3269 .list_runs(RunFilter::default(), 1, 10)
3270 .await
3271 .unwrap();
3272 assert_eq!(page.total, 0, "no run may be created");
3273 }
3274
3275 #[tokio::test]
3276 async fn enqueue_priority_bounds_are_accepted() {
3277 let mut engine = create_test_engine();
3278 engine.register(EchoWorkflow).unwrap();
3279
3280 for priority in [MIN_PRIORITY, MAX_PRIORITY] {
3281 let run = engine
3282 .enqueue_handler_with_options(
3283 "echo-workflow",
3284 TriggerKind::Api,
3285 json!({}),
3286 EnqueueOptions {
3287 priority: Some(priority),
3288 ..Default::default()
3289 },
3290 )
3291 .await
3292 .unwrap()
3293 .into_run();
3294 assert_eq!(run.priority, priority);
3295 }
3296 }
3297
3298 struct GpuWorkflow;
3299
3300 impl WorkflowHandler for GpuWorkflow {
3301 fn name(&self) -> &str {
3302 "gpu-workflow"
3303 }
3304
3305 fn required_worker_tags(&self) -> Vec<String> {
3306 vec!["gpu".to_string()]
3307 }
3308
3309 fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3310 Box::pin(async move { Ok(()) })
3311 }
3312 }
3313
3314 #[tokio::test]
3315 async fn enqueue_merges_handler_and_request_worker_tags() {
3316 let mut engine = create_test_engine();
3317 engine.register(GpuWorkflow).unwrap();
3318
3319 let run = engine
3320 .enqueue_handler_with_options(
3321 "gpu-workflow",
3322 TriggerKind::Api,
3323 json!({}),
3324 EnqueueOptions {
3325 worker_tags: vec!["region:eu".to_string(), "gpu".to_string()],
3326 ..Default::default()
3327 },
3328 )
3329 .await
3330 .unwrap()
3331 .into_run();
3332
3333 assert_eq!(
3334 run.worker_tags,
3335 vec!["gpu".to_string(), "region:eu".to_string()]
3336 );
3337 }
3338
3339 #[tokio::test]
3340 async fn enqueue_without_worker_tags_keeps_handler_tags() {
3341 let mut engine = create_test_engine();
3342 engine.register(GpuWorkflow).unwrap();
3343 engine.register(EchoWorkflow).unwrap();
3344
3345 let gpu = engine
3346 .enqueue_handler("gpu-workflow", TriggerKind::Api, json!({}), 0)
3347 .await
3348 .unwrap();
3349 assert_eq!(gpu.worker_tags, vec!["gpu".to_string()]);
3350
3351 let echo = engine
3352 .enqueue_handler("echo-workflow", TriggerKind::Api, json!({}), 0)
3353 .await
3354 .unwrap();
3355 assert!(echo.worker_tags.is_empty());
3356 }
3357
3358 #[tokio::test]
3359 async fn enqueue_rejects_invalid_worker_tags() {
3360 let mut engine = create_test_engine();
3361 engine.register(EchoWorkflow).unwrap();
3362
3363 for worker_tags in [
3364 vec!["bad,tag".to_string()],
3365 vec![" ".to_string()],
3366 vec!["x".repeat(65)],
3367 ] {
3368 let err = engine
3369 .enqueue_handler_with_options(
3370 "echo-workflow",
3371 TriggerKind::Api,
3372 json!({}),
3373 EnqueueOptions {
3374 worker_tags,
3375 ..Default::default()
3376 },
3377 )
3378 .await
3379 .unwrap_err();
3380 assert!(matches!(err, EngineError::InvalidWorkerTag(_)), "{err:?}");
3381 }
3382
3383 let err = engine
3385 .enqueue_handler_with_options(
3386 "not-registered",
3387 TriggerKind::Api,
3388 json!({}),
3389 EnqueueOptions {
3390 worker_tags: vec!["bad,tag".to_string()],
3391 ..Default::default()
3392 },
3393 )
3394 .await
3395 .unwrap_err();
3396 assert!(matches!(err, EngineError::InvalidWorkerTag(_)), "{err:?}");
3397 }
3398
3399 #[test]
3400 fn worker_tags_are_unset_by_default() {
3401 let engine = create_test_engine();
3402 assert!(engine.worker_tags().is_none());
3403 }
3404
3405 #[test]
3406 fn set_worker_tags_stores_the_tags() {
3407 let mut engine = create_test_engine();
3408 engine.set_worker_tags(vec!["gpu".to_string()]);
3409 assert_eq!(engine.worker_tags(), Some(&["gpu".to_string()][..]));
3410
3411 engine.set_worker_tags(Vec::new());
3412 assert_eq!(engine.worker_tags(), Some(&[][..]));
3413 }
3414
3415 #[tokio::test]
3416 async fn run_handler_records_handler_worker_tags() {
3417 let mut engine = create_test_engine();
3418 engine.register(GpuWorkflow).unwrap();
3419
3420 let result = engine
3421 .run_handler("gpu-workflow", TriggerKind::Manual, json!({}))
3422 .await
3423 .unwrap();
3424 assert_eq!(result.run.worker_tags, vec!["gpu".to_string()]);
3425 }
3426
3427 #[tokio::test]
3428 async fn run_handler_leaves_the_run_unattributed() {
3429 let mut engine = create_test_engine();
3430 engine.register(EchoWorkflow).unwrap();
3431
3432 let run = engine
3433 .run_handler("echo-workflow", TriggerKind::Manual, json!({}))
3434 .await
3435 .unwrap()
3436 .run;
3437
3438 assert!(run.created_by.is_none());
3439 }
3440
3441 #[tokio::test]
3442 async fn run_handler_priority_comes_from_the_handler() {
3443 let mut engine = create_test_engine();
3444 engine.register(UrgentWorkflow).unwrap();
3445
3446 let run = engine
3447 .run_handler("urgent-workflow", TriggerKind::Manual, json!({}))
3448 .await
3449 .unwrap()
3450 .run;
3451
3452 assert_eq!(run.priority, 60);
3453 }
3454
3455 #[tokio::test]
3456 async fn engine_register_boxed() {
3457 let mut engine = create_test_engine();
3458 let handler: Box<dyn WorkflowHandler> = Box::new(EchoWorkflow);
3459 let result = engine.register_boxed(handler);
3460 assert!(result.is_ok());
3461 assert_eq!(engine.handler_names().len(), 1);
3462 }
3463
3464 #[tokio::test]
3465 async fn engine_store_and_provider_accessors() {
3466 let store = Arc::new(InMemoryStore::new());
3467 let inner = ClaudeCodeProvider::new();
3468 let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
3469 inner,
3470 "/tmp/ironflow-fixtures",
3471 ));
3472 let engine = Engine::new(store.clone(), provider.clone());
3473
3474 let _ = engine.store();
3476 let _ = engine.provider();
3477 }
3478
3479 use crate::operation::{Operation, OperationContext};
3484 use async_trait::async_trait;
3485 use ironflow_core::error::OperationError;
3486 use ironflow_store::models::StepKind;
3487
3488 struct FakeGitlabOp {
3489 project_id: u64,
3490 title: String,
3491 }
3492
3493 #[async_trait]
3494 impl Operation for FakeGitlabOp {
3495 fn kind(&self) -> &str {
3496 "gitlab"
3497 }
3498
3499 async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
3500 Ok(json!({
3501 "issue_id": 42,
3502 "project_id": self.project_id,
3503 "title": self.title,
3504 }))
3505 }
3506
3507 fn input(&self) -> Option<Value> {
3508 Some(json!({
3509 "project_id": self.project_id,
3510 "title": self.title,
3511 }))
3512 }
3513 }
3514
3515 struct FailingOp;
3516
3517 #[async_trait]
3518 impl Operation for FailingOp {
3519 fn kind(&self) -> &str {
3520 "broken-service"
3521 }
3522
3523 async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
3524 Err(OperationError::Http {
3525 status: None,
3526 message: "service unavailable".to_string(),
3527 })
3528 }
3529 }
3530
3531 struct OperationWorkflow;
3532
3533 impl WorkflowHandler for OperationWorkflow {
3534 fn name(&self) -> &str {
3535 "operation-workflow"
3536 }
3537
3538 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3539 Box::pin(async move {
3540 let op = FakeGitlabOp {
3541 project_id: 123,
3542 title: "Bug report".to_string(),
3543 };
3544 ctx.operation("create-issue", &op).await?;
3545 Ok(())
3546 })
3547 }
3548 }
3549
3550 struct FailingOperationWorkflow;
3551
3552 impl WorkflowHandler for FailingOperationWorkflow {
3553 fn name(&self) -> &str {
3554 "failing-operation-workflow"
3555 }
3556
3557 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3558 Box::pin(async move {
3559 ctx.operation("broken-call", &FailingOp).await?;
3560 Ok(())
3561 })
3562 }
3563 }
3564
3565 struct MixedWorkflow;
3566
3567 impl WorkflowHandler for MixedWorkflow {
3568 fn name(&self) -> &str {
3569 "mixed-workflow"
3570 }
3571
3572 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3573 Box::pin(async move {
3574 ctx.shell("build", ShellConfig::new("echo built")).await?;
3575 let op = FakeGitlabOp {
3576 project_id: 456,
3577 title: "Deploy done".to_string(),
3578 };
3579 let result = ctx.operation("notify-gitlab", &op).await?;
3580 assert_eq!(result.output["issue_id"], 42);
3581 Ok(())
3582 })
3583 }
3584 }
3585
3586 #[tokio::test]
3587 async fn operation_step_happy_path() {
3588 let mut engine = create_test_engine();
3589 engine.register(OperationWorkflow).unwrap();
3590
3591 let run = engine
3592 .run_handler("operation-workflow", TriggerKind::Manual, json!({}))
3593 .await
3594 .unwrap()
3595 .run;
3596
3597 assert_eq!(run.status.state, RunStatus::Completed);
3598
3599 let steps = engine.store().list_steps(run.id).await.unwrap();
3600
3601 assert_eq!(steps.len(), 1);
3602 assert_eq!(steps[0].name, "create-issue");
3603 assert_eq!(steps[0].kind, StepKind::Custom("gitlab".to_string()));
3604 assert_eq!(
3605 steps[0].status.state,
3606 ironflow_store::models::StepStatus::Completed
3607 );
3608
3609 let output = steps[0].output.as_ref().unwrap();
3610 assert_eq!(output["issue_id"], 42);
3611 assert_eq!(output["project_id"], 123);
3612
3613 let input = steps[0].input.as_ref().unwrap();
3614 assert_eq!(input["project_id"], 123);
3615 assert_eq!(input["title"], "Bug report");
3616 }
3617
3618 #[tokio::test]
3619 async fn operation_step_failure_marks_run_failed() {
3620 let mut engine = create_test_engine();
3621 engine.register(FailingOperationWorkflow).unwrap();
3622
3623 let result = engine
3624 .run_handler("failing-operation-workflow", TriggerKind::Manual, json!({}))
3625 .await;
3626
3627 assert!(result.is_err());
3628 }
3629
3630 #[tokio::test]
3631 async fn operation_mixed_with_shell_steps() {
3632 let mut engine = create_test_engine();
3633 engine.register(MixedWorkflow).unwrap();
3634
3635 let run = engine
3636 .run_handler("mixed-workflow", TriggerKind::Manual, json!({}))
3637 .await
3638 .unwrap()
3639 .run;
3640
3641 assert_eq!(run.status.state, RunStatus::Completed);
3642
3643 let steps = engine.store().list_steps(run.id).await.unwrap();
3644
3645 assert_eq!(steps.len(), 2);
3646 assert_eq!(steps[0].kind, StepKind::Shell);
3647 assert_eq!(steps[1].kind, StepKind::Custom("gitlab".to_string()));
3648 assert_eq!(steps[0].position, 0);
3649 assert_eq!(steps[1].position, 1);
3650 }
3651
3652 use crate::config::ApprovalConfig;
3657
3658 struct SingleApprovalWorkflow;
3659
3660 impl WorkflowHandler for SingleApprovalWorkflow {
3661 fn name(&self) -> &str {
3662 "single-approval"
3663 }
3664
3665 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3666 Box::pin(async move {
3667 ctx.shell("build", ShellConfig::new("echo built")).await?;
3668 ctx.approval("gate", ApprovalConfig::new("OK?")).await?;
3669 ctx.shell("deploy", ShellConfig::new("echo deployed"))
3670 .await?;
3671 Ok(())
3672 })
3673 }
3674 }
3675
3676 struct DoubleApprovalWorkflow;
3677
3678 impl WorkflowHandler for DoubleApprovalWorkflow {
3679 fn name(&self) -> &str {
3680 "double-approval"
3681 }
3682
3683 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3684 Box::pin(async move {
3685 ctx.shell("build", ShellConfig::new("echo built")).await?;
3686 ctx.approval("staging-gate", ApprovalConfig::new("Deploy staging?"))
3687 .await?;
3688 ctx.shell("deploy-staging", ShellConfig::new("echo staging"))
3689 .await?;
3690 ctx.approval("prod-gate", ApprovalConfig::new("Deploy prod?"))
3691 .await?;
3692 ctx.shell("deploy-prod", ShellConfig::new("echo prod"))
3693 .await?;
3694 Ok(())
3695 })
3696 }
3697 }
3698
3699 #[tokio::test]
3700 async fn approval_pauses_run() {
3701 let mut engine = create_test_engine();
3702 engine.register(SingleApprovalWorkflow).unwrap();
3703
3704 let run = engine
3705 .run_handler("single-approval", TriggerKind::Manual, json!({}))
3706 .await
3707 .unwrap()
3708 .run;
3709
3710 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3711
3712 let steps = engine.store().list_steps(run.id).await.unwrap();
3713 assert_eq!(steps.len(), 2); assert_eq!(steps[0].kind, StepKind::Shell);
3715 assert_eq!(steps[0].status.state, StepStatus::Completed);
3716 assert_eq!(steps[1].kind, StepKind::Approval);
3717 assert_eq!(steps[1].status.state, StepStatus::AwaitingApproval);
3718 }
3719
3720 #[tokio::test]
3721 async fn approval_resume_completes_run() {
3722 let mut engine = create_test_engine();
3723 engine.register(SingleApprovalWorkflow).unwrap();
3724
3725 let run = engine
3727 .run_handler("single-approval", TriggerKind::Manual, json!({}))
3728 .await
3729 .unwrap()
3730 .run;
3731 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3732
3733 engine
3735 .store()
3736 .update_run_status(run.id, RunStatus::Running)
3737 .await
3738 .unwrap();
3739
3740 let resumed = engine.resume_run(run.id).await.unwrap().run;
3742 assert_eq!(resumed.status.state, RunStatus::Completed);
3743
3744 let steps = engine.store().list_steps(run.id).await.unwrap();
3745 assert_eq!(steps.len(), 3); assert_eq!(steps[0].name, "build");
3747 assert_eq!(steps[0].status.state, StepStatus::Completed);
3748 assert_eq!(steps[1].name, "gate");
3749 assert_eq!(steps[1].kind, StepKind::Approval);
3750 assert_eq!(steps[1].status.state, StepStatus::Completed);
3751 assert_eq!(steps[2].name, "deploy");
3752 assert_eq!(steps[2].status.state, StepStatus::Completed);
3753 }
3754
3755 #[tokio::test]
3756 async fn double_approval_two_resumes() {
3757 let mut engine = create_test_engine();
3758 engine.register(DoubleApprovalWorkflow).unwrap();
3759
3760 let run = engine
3762 .run_handler("double-approval", TriggerKind::Manual, json!({}))
3763 .await
3764 .unwrap()
3765 .run;
3766 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3767
3768 let steps = engine.store().list_steps(run.id).await.unwrap();
3769 assert_eq!(steps.len(), 2); engine
3773 .store()
3774 .update_run_status(run.id, RunStatus::Running)
3775 .await
3776 .unwrap();
3777
3778 let resumed = engine.resume_run(run.id).await.unwrap().run;
3779 assert_eq!(resumed.status.state, RunStatus::AwaitingApproval);
3780
3781 let steps = engine.store().list_steps(run.id).await.unwrap();
3782 assert_eq!(steps.len(), 4); engine
3786 .store()
3787 .update_run_status(run.id, RunStatus::Running)
3788 .await
3789 .unwrap();
3790
3791 let final_run = engine.resume_run(run.id).await.unwrap().run;
3792 assert_eq!(final_run.status.state, RunStatus::Completed);
3793
3794 let steps = engine.store().list_steps(run.id).await.unwrap();
3795 assert_eq!(steps.len(), 5);
3796 assert_eq!(steps[0].name, "build");
3797 assert_eq!(steps[1].name, "staging-gate");
3798 assert_eq!(steps[2].name, "deploy-staging");
3799 assert_eq!(steps[3].name, "prod-gate");
3800 assert_eq!(steps[4].name, "deploy-prod");
3801
3802 for step in &steps {
3803 assert_eq!(step.status.state, StepStatus::Completed);
3804 }
3805 }
3806
3807 use ironflow_store::models::{NewStep, StepUpdate, step_trace_id};
3812
3813 async fn create_step_with_status(
3814 store: &Arc<dyn Store>,
3815 run_id: Uuid,
3816 name: &str,
3817 position: u32,
3818 status: StepStatus,
3819 ) -> ironflow_store::models::Step {
3820 let step = store
3821 .create_step(NewStep {
3822 run_id,
3823 trace_id: step_trace_id(run_id, name, position),
3824 name: name.to_string(),
3825 kind: StepKind::Shell,
3826 position,
3827 input: None,
3828 is_error_handler: false,
3829 })
3830 .await
3831 .unwrap();
3832
3833 match status {
3834 StepStatus::Pending => {}
3835 StepStatus::Running => {
3836 store
3837 .update_step(
3838 step.id,
3839 StepUpdate {
3840 status: Some(StepStatus::Running),
3841 ..StepUpdate::default()
3842 },
3843 )
3844 .await
3845 .unwrap();
3846 }
3847 StepStatus::Completed => {
3848 store
3849 .update_step(
3850 step.id,
3851 StepUpdate {
3852 status: Some(StepStatus::Running),
3853 ..StepUpdate::default()
3854 },
3855 )
3856 .await
3857 .unwrap();
3858 store
3859 .update_step(
3860 step.id,
3861 StepUpdate {
3862 status: Some(StepStatus::Completed),
3863 ..StepUpdate::default()
3864 },
3865 )
3866 .await
3867 .unwrap();
3868 }
3869 StepStatus::AwaitingApproval => {
3870 store
3871 .update_step(
3872 step.id,
3873 StepUpdate {
3874 status: Some(StepStatus::Running),
3875 ..StepUpdate::default()
3876 },
3877 )
3878 .await
3879 .unwrap();
3880 store
3881 .update_step(
3882 step.id,
3883 StepUpdate {
3884 status: Some(StepStatus::AwaitingApproval),
3885 ..StepUpdate::default()
3886 },
3887 )
3888 .await
3889 .unwrap();
3890 }
3891 _ => panic!("unsupported status for test helper: {status}"),
3892 }
3893
3894 store.get_step(step.id).await.unwrap().unwrap()
3895 }
3896
3897 #[tokio::test]
3898 async fn fail_orphaned_steps_marks_running_as_failed() {
3899 let engine = create_test_engine();
3900 let run = engine
3901 .store()
3902 .create_run(NewRun {
3903 created_by: None,
3904 workflow_name: "test".to_string(),
3905 trigger: TriggerKind::Manual,
3906 payload: json!({}),
3907 max_retries: 0,
3908 handler_version: None,
3909 labels: HashMap::new(),
3910 scheduled_at: None,
3911 idempotency_key: None,
3912 concurrency_key: None,
3913 priority: 0,
3914 concurrency_limits: Vec::new(),
3915 max_cost_usd: None,
3916 worker_tags: Vec::new(),
3917 })
3918 .await
3919 .unwrap()
3920 .into_run();
3921
3922 let step = create_step_with_status(
3923 engine.store(),
3924 run.id,
3925 "running-step",
3926 0,
3927 StepStatus::Running,
3928 )
3929 .await;
3930
3931 engine
3932 .fail_orphaned_steps(run.id, "parent run timed out")
3933 .await
3934 .unwrap();
3935
3936 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
3937 assert_eq!(updated.status.state, StepStatus::Failed);
3938 assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
3939 assert!(updated.completed_at.is_some());
3940 }
3941
3942 #[tokio::test]
3943 async fn fail_orphaned_steps_marks_pending_as_skipped() {
3944 let engine = create_test_engine();
3945 let run = engine
3946 .store()
3947 .create_run(NewRun {
3948 created_by: None,
3949 workflow_name: "test".to_string(),
3950 trigger: TriggerKind::Manual,
3951 payload: json!({}),
3952 max_retries: 0,
3953 handler_version: None,
3954 labels: HashMap::new(),
3955 scheduled_at: None,
3956 idempotency_key: None,
3957 concurrency_key: None,
3958 priority: 0,
3959 concurrency_limits: Vec::new(),
3960 max_cost_usd: None,
3961 worker_tags: Vec::new(),
3962 })
3963 .await
3964 .unwrap()
3965 .into_run();
3966
3967 let step = create_step_with_status(
3968 engine.store(),
3969 run.id,
3970 "pending-step",
3971 0,
3972 StepStatus::Pending,
3973 )
3974 .await;
3975
3976 engine
3977 .fail_orphaned_steps(run.id, "parent run timed out")
3978 .await
3979 .unwrap();
3980
3981 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
3982 assert_eq!(updated.status.state, StepStatus::Skipped);
3983 assert!(updated.error.is_none());
3984 assert!(updated.completed_at.is_some());
3985 }
3986
3987 #[tokio::test]
3988 async fn fail_orphaned_steps_marks_awaiting_approval_as_failed() {
3989 let engine = create_test_engine();
3990 let run = engine
3991 .store()
3992 .create_run(NewRun {
3993 created_by: None,
3994 workflow_name: "test".to_string(),
3995 trigger: TriggerKind::Manual,
3996 payload: json!({}),
3997 max_retries: 0,
3998 handler_version: None,
3999 labels: HashMap::new(),
4000 scheduled_at: None,
4001 idempotency_key: None,
4002 concurrency_key: None,
4003 priority: 0,
4004 concurrency_limits: Vec::new(),
4005 max_cost_usd: None,
4006 worker_tags: Vec::new(),
4007 })
4008 .await
4009 .unwrap()
4010 .into_run();
4011
4012 let step = create_step_with_status(
4013 engine.store(),
4014 run.id,
4015 "approval-step",
4016 0,
4017 StepStatus::AwaitingApproval,
4018 )
4019 .await;
4020
4021 engine
4022 .fail_orphaned_steps(run.id, "parent run timed out")
4023 .await
4024 .unwrap();
4025
4026 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
4027 assert_eq!(updated.status.state, StepStatus::Failed);
4028 assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
4029 assert!(updated.completed_at.is_some());
4030 }
4031
4032 #[tokio::test]
4033 async fn fail_orphaned_steps_skips_terminal_steps() {
4034 let engine = create_test_engine();
4035 let run = engine
4036 .store()
4037 .create_run(NewRun {
4038 created_by: None,
4039 workflow_name: "test".to_string(),
4040 trigger: TriggerKind::Manual,
4041 payload: json!({}),
4042 max_retries: 0,
4043 handler_version: None,
4044 labels: HashMap::new(),
4045 scheduled_at: None,
4046 idempotency_key: None,
4047 concurrency_key: None,
4048 priority: 0,
4049 concurrency_limits: Vec::new(),
4050 max_cost_usd: None,
4051 worker_tags: Vec::new(),
4052 })
4053 .await
4054 .unwrap()
4055 .into_run();
4056
4057 let completed_step =
4058 create_step_with_status(engine.store(), run.id, "done", 0, StepStatus::Completed).await;
4059 let running_step =
4060 create_step_with_status(engine.store(), run.id, "in-flight", 1, StepStatus::Running)
4061 .await;
4062
4063 engine
4064 .fail_orphaned_steps(run.id, "parent run timed out")
4065 .await
4066 .unwrap();
4067
4068 let completed = engine
4069 .store()
4070 .get_step(completed_step.id)
4071 .await
4072 .unwrap()
4073 .unwrap();
4074 assert_eq!(completed.status.state, StepStatus::Completed);
4075
4076 let failed = engine
4077 .store()
4078 .get_step(running_step.id)
4079 .await
4080 .unwrap()
4081 .unwrap();
4082 assert_eq!(failed.status.state, StepStatus::Failed);
4083 }
4084
4085 #[tokio::test]
4086 async fn fail_orphaned_steps_mixed_states() {
4087 let engine = create_test_engine();
4088 let run = engine
4089 .store()
4090 .create_run(NewRun {
4091 created_by: None,
4092 workflow_name: "test".to_string(),
4093 trigger: TriggerKind::Manual,
4094 payload: json!({}),
4095 max_retries: 0,
4096 handler_version: None,
4097 labels: HashMap::new(),
4098 scheduled_at: None,
4099 idempotency_key: None,
4100 concurrency_key: None,
4101 priority: 0,
4102 concurrency_limits: Vec::new(),
4103 max_cost_usd: None,
4104 worker_tags: Vec::new(),
4105 })
4106 .await
4107 .unwrap()
4108 .into_run();
4109
4110 let s_completed =
4111 create_step_with_status(engine.store(), run.id, "step-1", 0, StepStatus::Completed)
4112 .await;
4113 let s_running =
4114 create_step_with_status(engine.store(), run.id, "step-2", 1, StepStatus::Running).await;
4115 let s_pending =
4116 create_step_with_status(engine.store(), run.id, "step-3", 2, StepStatus::Pending).await;
4117
4118 engine.fail_orphaned_steps(run.id, "timeout").await.unwrap();
4119
4120 let r_completed = engine
4121 .store()
4122 .get_step(s_completed.id)
4123 .await
4124 .unwrap()
4125 .unwrap();
4126 assert_eq!(r_completed.status.state, StepStatus::Completed);
4127
4128 let r_running = engine
4129 .store()
4130 .get_step(s_running.id)
4131 .await
4132 .unwrap()
4133 .unwrap();
4134 assert_eq!(r_running.status.state, StepStatus::Failed);
4135 assert_eq!(r_running.error.as_deref(), Some("timeout"));
4136
4137 let r_pending = engine
4138 .store()
4139 .get_step(s_pending.id)
4140 .await
4141 .unwrap()
4142 .unwrap();
4143 assert_eq!(r_pending.status.state, StepStatus::Skipped);
4144 assert!(r_pending.error.is_none());
4145 }
4146
4147 #[tokio::test]
4148 async fn fail_orphaned_steps_no_steps_is_noop() {
4149 let engine = create_test_engine();
4150 let run = engine
4151 .store()
4152 .create_run(NewRun {
4153 created_by: None,
4154 workflow_name: "test".to_string(),
4155 trigger: TriggerKind::Manual,
4156 payload: json!({}),
4157 max_retries: 0,
4158 handler_version: None,
4159 labels: HashMap::new(),
4160 scheduled_at: None,
4161 idempotency_key: None,
4162 concurrency_key: None,
4163 priority: 0,
4164 concurrency_limits: Vec::new(),
4165 max_cost_usd: None,
4166 worker_tags: Vec::new(),
4167 })
4168 .await
4169 .unwrap()
4170 .into_run();
4171
4172 let result = engine.fail_orphaned_steps(run.id, "timeout").await;
4173 assert!(result.is_ok());
4174 }
4175
4176 #[tokio::test]
4177 async fn fail_orphaned_steps_preserves_existing_error() {
4178 let engine = create_test_engine();
4179 let run = engine
4180 .store()
4181 .create_run(NewRun {
4182 created_by: None,
4183 workflow_name: "test".to_string(),
4184 trigger: TriggerKind::Manual,
4185 payload: json!({}),
4186 max_retries: 0,
4187 handler_version: None,
4188 labels: HashMap::new(),
4189 scheduled_at: None,
4190 idempotency_key: None,
4191 concurrency_key: None,
4192 priority: 0,
4193 concurrency_limits: Vec::new(),
4194 max_cost_usd: None,
4195 worker_tags: Vec::new(),
4196 })
4197 .await
4198 .unwrap()
4199 .into_run();
4200
4201 let step_with_error = create_step_with_status(
4202 engine.store(),
4203 run.id,
4204 "already-errored",
4205 0,
4206 StepStatus::Running,
4207 )
4208 .await;
4209
4210 engine
4211 .store()
4212 .update_step(
4213 step_with_error.id,
4214 StepUpdate {
4215 error: Some("real error from provider".to_string()),
4216 ..StepUpdate::default()
4217 },
4218 )
4219 .await
4220 .unwrap();
4221
4222 let step_no_error = create_step_with_status(
4223 engine.store(),
4224 run.id,
4225 "no-error-yet",
4226 1,
4227 StepStatus::Running,
4228 )
4229 .await;
4230
4231 engine
4232 .fail_orphaned_steps(run.id, "parent run failed")
4233 .await
4234 .unwrap();
4235
4236 let updated_with = engine
4237 .store()
4238 .get_step(step_with_error.id)
4239 .await
4240 .unwrap()
4241 .unwrap();
4242 assert_eq!(updated_with.status.state, StepStatus::Failed);
4243 assert_eq!(
4244 updated_with.error.as_deref(),
4245 Some("real error from provider"),
4246 );
4247
4248 let updated_without = engine
4249 .store()
4250 .get_step(step_no_error.id)
4251 .await
4252 .unwrap()
4253 .unwrap();
4254 assert_eq!(updated_without.status.state, StepStatus::Failed);
4255 assert_eq!(updated_without.error.as_deref(), Some("parent run failed"),);
4256 }
4257}