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(
1846 &self,
1847 run_id: Uuid,
1848 error: &str,
1849 retryable: bool,
1850 cost_usd: Option<Decimal>,
1851 duration_ms: Option<u64>,
1852 ) -> Result<RunStatus, EngineError> {
1853 let run = self
1854 .store
1855 .get_run(run_id)
1856 .await?
1857 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
1858
1859 let has_attempts_left = run.retry_count < run.max_retries;
1860 let update = if retryable && has_attempts_left {
1861 let backoff = backoff_for_retry(run.retry_count);
1862 let scheduled_at = Utc::now() + TimeDelta::milliseconds(backoff.as_millis() as i64);
1863
1864 info!(
1865 run_id = %run_id,
1866 workflow = %run.workflow_name,
1867 attempt = run.retry_count + 1,
1868 max_retries = run.max_retries,
1869 backoff_secs = backoff.as_secs(),
1870 scheduled_at = %scheduled_at,
1871 "run failed, scheduling retry"
1872 );
1873
1874 RunUpdate {
1875 status: Some(RunStatus::Retrying),
1876 error: Some(error.to_string()),
1877 increment_retry: true,
1878 cost_usd,
1879 duration_ms,
1880 scheduled_at: Some(scheduled_at),
1881 ..RunUpdate::default()
1882 }
1883 } else {
1884 RunUpdate {
1885 status: Some(RunStatus::Failed),
1886 error: Some(error.to_string()),
1887 cost_usd,
1888 duration_ms,
1889 completed_at: Some(Utc::now()),
1890 ..RunUpdate::default()
1891 }
1892 };
1893
1894 let status = update.status.unwrap_or(RunStatus::Failed);
1895 self.store.update_run(run_id, update).await?;
1896 self.fail_orphaned_steps(run_id, error).await?;
1897 self.cancel_descendants_of_stopped_run(run_id, error).await;
1900
1901 Ok(status)
1902 }
1903
1904 pub async fn interrupt_running_steps(&self, run_id: Uuid) -> Result<(), EngineError> {
1934 interrupt_running_steps(self.store.as_ref(), run_id).await
1935 }
1936
1937 pub async fn fail_orphaned_steps(
1951 &self,
1952 run_id: Uuid,
1953 error_message: &str,
1954 ) -> Result<(), EngineError> {
1955 let steps = self.store.list_steps(run_id).await?;
1956 let now = Utc::now();
1957
1958 for step in steps {
1959 if step.status.state.is_terminal() {
1960 continue;
1961 }
1962
1963 let (target_status, error) = match step.status.state {
1964 StepStatus::Running | StepStatus::AwaitingApproval => {
1965 let err = if step.error.is_some() {
1966 None
1967 } else {
1968 Some(error_message.to_string())
1969 };
1970 (StepStatus::Failed, err)
1971 }
1972 StepStatus::Pending => (StepStatus::Skipped, None),
1973 _ => continue,
1974 };
1975
1976 if let Err(e) = self
1977 .store
1978 .update_step(
1979 step.id,
1980 StepUpdate {
1981 status: Some(target_status),
1982 error,
1983 completed_at: Some(now),
1984 ..StepUpdate::default()
1985 },
1986 )
1987 .await
1988 {
1989 warn!(
1990 run_id = %run_id,
1991 step_id = %step.id,
1992 step_name = %step.name,
1993 error = %e,
1994 "failed to cleanup orphaned step"
1995 );
1996 } else {
1997 info!(
1998 run_id = %run_id,
1999 step_id = %step.id,
2000 step_name = %step.name,
2001 from = %step.status.state,
2002 to = %target_status,
2003 "cleaned up orphaned step"
2004 );
2005 }
2006 }
2007
2008 Ok(())
2009 }
2010
2011 async fn release_then_execute(
2016 &self,
2017 run_id: Uuid,
2018 handler: &dyn WorkflowHandler,
2019 ctx: &mut WorkflowContext,
2020 ) -> Result<(), EngineError> {
2021 match self.provider.release_run(&run_id.to_string()).await {
2022 Ok(()) => handler.execute(ctx).await,
2023 Err(e) => Err(EngineError::Operation(OperationError::Agent(e))),
2024 }
2025 }
2026
2027 async fn finalize_run(
2033 &self,
2034 run_id: Uuid,
2035 workflow_name: &str,
2036 result: Result<(), EngineError>,
2037 ctx: &WorkflowContext,
2038 run_start: Instant,
2039 run_labels: HashMap<String, String>,
2040 ) -> Result<WorkflowResult, EngineError> {
2041 let total_duration = ctx.carried_duration_ms() + run_start.elapsed().as_millis() as u64;
2044 let completed_at = Utc::now();
2045
2046 let final_status;
2047 let final_run;
2048
2049 match result {
2050 Ok(()) => {
2051 final_status = if ctx.has_allowed_failure() {
2052 RunStatus::Warning
2053 } else {
2054 RunStatus::Completed
2055 };
2056 final_run = self
2057 .store
2058 .update_run_returning(
2059 run_id,
2060 RunUpdate {
2061 status: Some(final_status),
2062 cost_usd: Some(ctx.total_cost_usd()),
2063 duration_ms: Some(total_duration),
2064 completed_at: Some(completed_at),
2065 output: ctx.output().cloned(),
2066 ..RunUpdate::default()
2067 },
2068 )
2069 .await?;
2070
2071 info!(
2072 run_id = %run_id,
2073 status = %final_status,
2074 cost_usd = %ctx.total_cost_usd(),
2075 duration_ms = total_duration,
2076 "run completed"
2077 );
2078 }
2079 Err(EngineError::ApprovalRequired {
2080 run_id: approval_run_id,
2081 step_id,
2082 ref message,
2083 }) => {
2084 final_status = RunStatus::AwaitingApproval;
2085 final_run = self
2086 .store
2087 .update_run_returning(
2088 run_id,
2089 RunUpdate {
2090 status: Some(RunStatus::AwaitingApproval),
2091 cost_usd: Some(ctx.total_cost_usd()),
2092 duration_ms: Some(total_duration),
2093 ..RunUpdate::default()
2094 },
2095 )
2096 .await?;
2097
2098 info!(
2099 run_id = %approval_run_id,
2100 step_id = %step_id,
2101 message = %message,
2102 "run awaiting approval"
2103 );
2104
2105 self.publish_approval_requested(approval_run_id, step_id, message)
2106 .await?;
2107 }
2108 Err(EngineError::ChildSuspended {
2109 run_id: child_run_id,
2110 ref cause,
2111 }) => {
2112 final_status = cause.suspension_status();
2113 final_run = self
2117 .store
2118 .update_run_returning(
2119 run_id,
2120 RunUpdate {
2121 status: Some(final_status),
2122 cost_usd: Some(ctx.total_cost_usd()),
2123 duration_ms: Some(total_duration),
2124 ..RunUpdate::default()
2125 },
2126 )
2127 .await?;
2128
2129 let leaf = cause.suspension_leaf();
2130 info!(
2131 run_id = %run_id,
2132 child_run_id = %child_run_id,
2133 status = %final_status,
2134 cause = %leaf,
2135 "run suspended with its child run"
2136 );
2137
2138 match leaf {
2139 EngineError::ApprovalRequired {
2140 run_id: approval_run_id,
2141 step_id,
2142 message,
2143 } => {
2144 self.publish_approval_requested(*approval_run_id, *step_id, message)
2145 .await?;
2146 }
2147 EngineError::SignalWaiting {
2148 run_id: wait_run_id,
2149 step_id,
2150 step_name,
2151 name,
2152 key,
2153 deadline_at,
2154 } => {
2155 self.event_publisher
2156 .publish(Event::SignalAwaited(SignalAwaitedEvent {
2157 run_id: *wait_run_id,
2158 step_id: *step_id,
2159 step_name: step_name.clone(),
2160 name: name.clone(),
2161 key: key.clone(),
2162 deadline_at: *deadline_at,
2163 at: Utc::now(),
2164 }));
2165 }
2166 _ => {}
2169 }
2170 }
2171 Err(EngineError::HumanInputRequired {
2172 run_id: input_run_id,
2173 step_id,
2174 ref message,
2175 }) => {
2176 final_status = RunStatus::AwaitingApproval;
2177 final_run = self
2178 .store
2179 .update_run_returning(
2180 run_id,
2181 RunUpdate {
2182 status: Some(RunStatus::AwaitingApproval),
2183 cost_usd: Some(ctx.total_cost_usd()),
2184 duration_ms: Some(total_duration),
2185 ..RunUpdate::default()
2186 },
2187 )
2188 .await?;
2189
2190 info!(
2192 run_id = %input_run_id,
2193 step_id = %step_id,
2194 message = %message,
2195 "run awaiting human input"
2196 );
2197 }
2198 Err(EngineError::DelaySleeping {
2199 run_id: delay_run_id,
2200 step_id,
2201 wake_at,
2202 }) => {
2203 final_status = RunStatus::Sleeping;
2204 final_run = self
2205 .store
2206 .update_run_returning(
2207 run_id,
2208 RunUpdate {
2209 status: Some(RunStatus::Sleeping),
2210 cost_usd: Some(ctx.total_cost_usd()),
2211 duration_ms: Some(total_duration),
2212 scheduled_at: Some(wake_at),
2213 ..RunUpdate::default()
2214 },
2215 )
2216 .await?;
2217
2218 info!(
2219 run_id = %delay_run_id,
2220 step_id = %step_id,
2221 wake_at = %wake_at,
2222 "run sleeping until delay elapses"
2223 );
2224 }
2225 Err(EngineError::CapacitySleeping {
2226 run_id: capacity_run_id,
2227 step_id,
2228 ref kind,
2229 wake_at,
2230 }) => {
2231 final_status = RunStatus::Sleeping;
2232 final_run = self
2233 .store
2234 .update_run_returning(
2235 run_id,
2236 RunUpdate {
2237 status: Some(RunStatus::Sleeping),
2238 cost_usd: Some(ctx.total_cost_usd()),
2239 duration_ms: Some(total_duration),
2240 scheduled_at: Some(wake_at),
2241 capacity_wait_kind: Some(ProviderKind::new(kind.as_str())),
2242 ..RunUpdate::default()
2243 },
2244 )
2245 .await?;
2246
2247 info!(
2248 run_id = %capacity_run_id,
2249 step_id = %step_id,
2250 kind = %kind,
2251 wake_at = %wake_at,
2252 "run sleeping until provider capacity returns"
2253 );
2254 }
2255 Err(EngineError::SignalWaiting {
2256 run_id: wait_run_id,
2257 step_id,
2258 ref step_name,
2259 ref name,
2260 ref key,
2261 deadline_at,
2262 }) => {
2263 final_status = RunStatus::Sleeping;
2264 let waiting = self
2268 .store
2269 .suspend_run_on_signal(run_id, step_id, deadline_at)
2270 .await?;
2271 final_run = self
2272 .store
2273 .update_run_returning(
2274 run_id,
2275 RunUpdate {
2276 cost_usd: Some(ctx.total_cost_usd()),
2277 duration_ms: Some(total_duration),
2278 ..RunUpdate::default()
2279 },
2280 )
2281 .await?;
2282
2283 if waiting {
2284 self.event_publisher
2285 .publish(Event::SignalAwaited(SignalAwaitedEvent {
2286 run_id: wait_run_id,
2287 step_id,
2288 step_name: step_name.clone(),
2289 name: name.clone(),
2290 key: key.clone(),
2291 deadline_at,
2292 at: Utc::now(),
2293 }));
2294 }
2295
2296 info!(
2297 run_id = %wait_run_id,
2298 step_id = %step_id,
2299 signal = %name,
2300 key = %key,
2301 deadline_at = %deadline_at,
2302 waiting,
2303 "run sleeping until a signal arrives"
2304 );
2305 }
2306 Err(err) => {
2307 let guardrail_stop = matches!(
2311 err,
2312 EngineError::RunBudgetExceeded { .. } | EngineError::WorkflowGuardRejected(_)
2313 );
2314
2315 final_status = if guardrail_stop {
2316 if let Err(store_err) = self
2317 .store
2318 .update_run(
2319 run_id,
2320 RunUpdate {
2321 status: Some(RunStatus::Cancelled),
2322 error: Some(err.to_string()),
2323 cost_usd: Some(ctx.total_cost_usd()),
2324 duration_ms: Some(total_duration),
2325 completed_at: Some(completed_at),
2326 output: ctx.output().cloned(),
2327 ..RunUpdate::default()
2328 },
2329 )
2330 .await
2331 {
2332 error!(run_id = %run_id, store_error = %store_err, "failed to persist run cancellation");
2333 }
2334 if let Err(cleanup_err) = self
2335 .fail_orphaned_steps(run_id, "run stopped: guardrail limit reached")
2336 .await
2337 {
2338 error!(run_id = %run_id, store_error = %cleanup_err, "failed to cleanup orphaned steps");
2339 }
2340 RunStatus::Cancelled
2341 } else {
2342 if let Some(output) = ctx.output()
2345 && let Err(store_err) = self
2346 .store
2347 .update_run(
2348 run_id,
2349 RunUpdate {
2350 output: Some(output.clone()),
2351 ..RunUpdate::default()
2352 },
2353 )
2354 .await
2355 {
2356 error!(run_id = %run_id, store_error = %store_err, "failed to persist run output");
2357 }
2358 self.fail_or_schedule_retry(
2359 run_id,
2360 &err.to_string(),
2361 is_run_retryable(&err),
2362 Some(ctx.total_cost_usd()),
2363 Some(total_duration),
2364 )
2365 .await
2366 .unwrap_or_else(|store_err| {
2367 error!(run_id = %run_id, store_error = %store_err, "failed to persist run failure");
2368 RunStatus::Failed
2369 })
2370 };
2371
2372 if matches!(err, EngineError::RunBudgetExceeded { .. }) {
2373 self.on_run_budget_exceeded(workflow_name, run_id, &err);
2374 }
2375
2376 error!(run_id = %run_id, status = %final_status, error = %err, "run stopped");
2377
2378 self.publish_run_status_changed(
2379 workflow_name,
2380 run_id,
2381 final_status,
2382 Some(err.to_string()),
2383 ctx,
2384 total_duration,
2385 run_labels,
2386 );
2387
2388 #[cfg(feature = "prometheus")]
2389 self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
2390
2391 return Err(err);
2392 }
2393 }
2394
2395 self.publish_run_status_changed(
2396 workflow_name,
2397 run_id,
2398 final_status,
2399 None,
2400 ctx,
2401 total_duration,
2402 run_labels,
2403 );
2404
2405 #[cfg(feature = "prometheus")]
2406 self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
2407
2408 Ok(WorkflowResult {
2409 run: final_run,
2410 steps: ctx.step_results().to_vec(),
2411 })
2412 }
2413
2414 async fn publish_approval_requested(
2417 &self,
2418 run_id: Uuid,
2419 step_id: Uuid,
2420 message: &str,
2421 ) -> Result<(), EngineError> {
2422 let requirement = self
2423 .store
2424 .get_step(step_id)
2425 .await?
2426 .and_then(|s| s.approval_requirement);
2427 self.event_publisher
2428 .publish(Event::ApprovalRequested(ApprovalRequestedEvent {
2429 run_id,
2430 step_id,
2431 message: message.to_string(),
2432 requirement,
2433 at: Utc::now(),
2434 }));
2435 Ok(())
2436 }
2437
2438 pub async fn fail_ancestors(&self, run_id: Uuid, reason: &str) -> Result<(), EngineError> {
2467 let mut current = self
2468 .store
2469 .get_run(run_id)
2470 .await?
2471 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
2472 let mut visited = HashSet::from([run_id]);
2474
2475 while let Some(parent_id) = chain_parent(¤t) {
2476 if !visited.insert(parent_id) {
2477 break;
2478 }
2479 let status = self
2480 .fail_or_schedule_retry(parent_id, reason, false, None, None)
2481 .await?;
2482 info!(
2483 run_id = %run_id,
2484 ancestor_run_id = %parent_id,
2485 status = %status,
2486 "ancestor run failed with its child"
2487 );
2488 current = self
2489 .store
2490 .get_run(parent_id)
2491 .await?
2492 .ok_or(EngineError::Store(StoreError::RunNotFound(parent_id)))?;
2493 }
2494
2495 Ok(())
2496 }
2497
2498 #[cfg(feature = "prometheus")]
2500 fn emit_run_metrics(
2501 &self,
2502 workflow_name: &str,
2503 status: RunStatus,
2504 duration_ms: u64,
2505 ctx: &WorkflowContext,
2506 ) {
2507 let status_str = status.to_string();
2508 let wf = workflow_name.to_string();
2509
2510 counter!(RUNS_TOTAL, "workflow" => wf.clone(), "status" => status_str.clone()).increment(1);
2511 histogram!(RUN_DURATION_SECONDS, "workflow" => wf.clone(), "status" => status_str)
2512 .record(duration_ms as f64 / 1000.0);
2513 histogram!(RUN_COST_USD, "workflow" => wf.clone()).record(
2514 ctx.total_cost_usd()
2515 .to_string()
2516 .parse::<f64>()
2517 .unwrap_or(0.0),
2518 );
2519 gauge!(RUNS_ACTIVE, "workflow" => wf).decrement(1.0);
2520 }
2521
2522 fn on_run_budget_exceeded(&self, workflow_name: &str, run_id: Uuid, err: &EngineError) {
2528 let EngineError::RunBudgetExceeded {
2529 limit_usd,
2530 spent_usd,
2531 step_budget_usd,
2532 ..
2533 } = err
2534 else {
2535 return;
2536 };
2537
2538 #[cfg(feature = "prometheus")]
2539 counter!(
2540 RUN_BUDGET_EXCEEDED_TOTAL,
2541 "workflow" => workflow_name.to_string(),
2542 "scope" => "run",
2543 )
2544 .increment(1);
2545
2546 self.event_publisher
2547 .publish(Event::RunBudgetExceeded(RunBudgetExceededEvent {
2548 run_id,
2549 workflow_name: workflow_name.to_string(),
2550 limit_usd: *limit_usd,
2551 spent_usd: *spent_usd,
2552 step_budget_usd: *step_budget_usd,
2553 at: Utc::now(),
2554 }));
2555 }
2556
2557 #[allow(clippy::too_many_arguments)]
2562 fn publish_run_status_changed(
2563 &self,
2564 workflow_name: &str,
2565 run_id: Uuid,
2566 to: RunStatus,
2567 error: Option<String>,
2568 ctx: &WorkflowContext,
2569 duration_ms: u64,
2570 labels: HashMap<String, String>,
2571 ) {
2572 let now = Utc::now();
2573 let cost_usd = ctx.total_cost_usd();
2574 let wf = workflow_name.to_string();
2575
2576 self.event_publisher
2577 .publish(Event::RunStatusChanged(RunStatusChangedEvent {
2578 run_id,
2579 workflow_name: wf.clone(),
2580 from: RunStatus::Running,
2581 to,
2582 error: error.clone(),
2583 cost_usd,
2584 duration_ms,
2585 labels: labels.clone(),
2586 at: now,
2587 }));
2588
2589 if to == RunStatus::Failed {
2590 self.event_publisher
2591 .publish(Event::RunFailed(RunFailedEvent {
2592 run_id,
2593 workflow_name: wf,
2594 error,
2595 cost_usd,
2596 duration_ms,
2597 labels,
2598 at: now,
2599 }));
2600 }
2601 }
2602}
2603
2604impl fmt::Debug for Engine {
2605 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
2606 f.debug_struct("Engine")
2607 .field("handlers", &self.handlers.keys().collect::<Vec<_>>())
2608 .finish_non_exhaustive()
2609 }
2610}
2611
2612#[cfg(test)]
2613mod tests {
2614 use super::*;
2615 use crate::config::ShellConfig;
2616 use crate::handler::{HandlerFuture, WorkflowHandler};
2617 use ironflow_core::providers::claude::ClaudeCodeProvider;
2618 use ironflow_core::providers::record_replay::RecordReplayProvider;
2619 use ironflow_store::memory::InMemoryStore;
2620 use ironflow_store::models::{MAX_PRIORITY, MIN_PRIORITY, StepStatus};
2621 use serde_json::json;
2622
2623 struct EchoWorkflow;
2625
2626 impl WorkflowHandler for EchoWorkflow {
2627 fn name(&self) -> &str {
2628 "echo-workflow"
2629 }
2630
2631 fn describe(&self) -> WorkflowInfo {
2632 WorkflowInfo {
2633 description: "A simple workflow that echoes hello".to_string(),
2634 source_code: None,
2635 sub_workflows: Vec::new(),
2636 category: None,
2637 version: self.version().map(str::to_string),
2638 compatible_versions: Vec::new(),
2639 input_schema: None,
2640 default_labels: HashMap::new(),
2641 schedule: self.schedule().cloned(),
2642 default_max_cost_usd: self.default_max_cost_usd(),
2643 priority: self.priority(),
2644 }
2645 }
2646
2647 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2648 Box::pin(async move {
2649 ctx.shell("greet", ShellConfig::new("echo hello")).await?;
2650 Ok(())
2651 })
2652 }
2653 }
2654
2655 struct FailingWorkflow;
2657
2658 impl WorkflowHandler for FailingWorkflow {
2659 fn name(&self) -> &str {
2660 "failing-workflow"
2661 }
2662
2663 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2664 Box::pin(async move {
2665 ctx.shell("fail", ShellConfig::new("exit 1")).await?;
2666 Ok(())
2667 })
2668 }
2669 }
2670
2671 fn create_test_engine() -> Engine {
2672 let store = Arc::new(InMemoryStore::new());
2673 let inner = ClaudeCodeProvider::new();
2674 let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
2675 inner,
2676 "/tmp/ironflow-fixtures",
2677 ));
2678 Engine::new(store, provider)
2679 }
2680
2681 #[test]
2682 fn engine_new_creates_instance() {
2683 let engine = create_test_engine();
2684 assert_eq!(engine.handler_names().len(), 0);
2685 }
2686
2687 #[test]
2688 fn execution_mode_defaults_to_local() {
2689 let engine = create_test_engine();
2690 assert_eq!(engine.execution_mode(), ExecutionMode::Local);
2691 }
2692
2693 #[test]
2694 fn with_execution_mode_overrides_the_default() {
2695 let engine = create_test_engine().with_execution_mode(ExecutionMode::Workers);
2696 assert_eq!(engine.execution_mode(), ExecutionMode::Workers);
2697 }
2698
2699 #[test]
2700 fn engine_register_handler() {
2701 let mut engine = create_test_engine();
2702 let result = engine.register(EchoWorkflow);
2703 assert!(result.is_ok());
2704 assert_eq!(engine.handler_names().len(), 1);
2705 assert!(engine.handler_names().contains(&"echo-workflow"));
2706 }
2707
2708 #[test]
2709 fn engine_register_duplicate_returns_error() {
2710 let mut engine = create_test_engine();
2711 engine.register(EchoWorkflow).unwrap();
2712 let result = engine.register(EchoWorkflow);
2713 assert!(result.is_err());
2714 }
2715
2716 #[test]
2717 fn engine_get_handler_found() {
2718 let mut engine = create_test_engine();
2719 engine.register(EchoWorkflow).unwrap();
2720 let handler = engine.get_handler("echo-workflow");
2721 assert!(handler.is_some());
2722 }
2723
2724 #[test]
2725 fn engine_get_handler_not_found() {
2726 let engine = create_test_engine();
2727 let handler = engine.get_handler("nonexistent");
2728 assert!(handler.is_none());
2729 }
2730
2731 #[test]
2732 fn engine_handler_names_lists_all() {
2733 let mut engine = create_test_engine();
2734 engine.register(EchoWorkflow).unwrap();
2735 engine.register(FailingWorkflow).unwrap();
2736 let names = engine.handler_names();
2737 assert_eq!(names.len(), 2);
2738 assert!(names.contains(&"echo-workflow"));
2739 assert!(names.contains(&"failing-workflow"));
2740 }
2741
2742 #[test]
2743 fn engine_handler_info_returns_description() {
2744 let mut engine = create_test_engine();
2745 engine.register(EchoWorkflow).unwrap();
2746 let info = engine.handler_info("echo-workflow");
2747 assert!(info.is_some());
2748 let info = info.unwrap();
2749 assert_eq!(info.description, "A simple workflow that echoes hello");
2750 }
2751
2752 struct CategorizedWorkflow;
2753
2754 impl WorkflowHandler for CategorizedWorkflow {
2755 fn name(&self) -> &str {
2756 "categorized"
2757 }
2758 fn category(&self) -> Option<&str> {
2759 Some("data/etl")
2760 }
2761 fn execute<'a>(
2762 &'a self,
2763 _ctx: &'a mut WorkflowContext,
2764 ) -> crate::handler::HandlerFuture<'a> {
2765 Box::pin(async move { Ok(()) })
2766 }
2767 }
2768
2769 #[test]
2770 fn engine_default_describe_propagates_category() {
2771 let mut engine = create_test_engine();
2772 engine.register(CategorizedWorkflow).unwrap();
2773 let info = engine.handler_info("categorized").unwrap();
2774 assert_eq!(info.category.as_deref(), Some("data/etl"));
2775 }
2776
2777 #[test]
2778 fn engine_default_describe_without_category() {
2779 let mut engine = create_test_engine();
2780 engine.register(EchoWorkflow).unwrap();
2781 let info = engine.handler_info("echo-workflow").unwrap();
2782 assert!(info.category.is_none());
2783 }
2784
2785 struct ScheduledWorkflow {
2790 schedule: CronSchedule,
2791 }
2792
2793 impl ScheduledWorkflow {
2794 fn new() -> Self {
2795 Self {
2796 schedule: CronSchedule::new("0 0 * * * *").unwrap(),
2797 }
2798 }
2799 }
2800
2801 impl WorkflowHandler for ScheduledWorkflow {
2802 fn name(&self) -> &str {
2803 "scheduled"
2804 }
2805 fn schedule(&self) -> Option<&CronSchedule> {
2806 Some(&self.schedule)
2807 }
2808 fn execute<'a>(
2809 &'a self,
2810 _ctx: &'a mut WorkflowContext,
2811 ) -> crate::handler::HandlerFuture<'a> {
2812 Box::pin(async move { Ok(()) })
2813 }
2814 }
2815
2816 #[test]
2817 fn engine_default_describe_propagates_schedule() {
2818 let mut engine = create_test_engine();
2819 engine.register(ScheduledWorkflow::new()).unwrap();
2820 let info = engine.handler_info("scheduled").unwrap();
2821 assert_eq!(
2822 info.schedule.as_ref().map(|s| s.as_str()),
2823 Some("0 0 * * * *")
2824 );
2825 }
2826
2827 #[test]
2828 fn engine_default_describe_without_schedule() {
2829 let mut engine = create_test_engine();
2830 engine.register(EchoWorkflow).unwrap();
2831 let info = engine.handler_info("echo-workflow").unwrap();
2832 assert!(info.schedule.is_none());
2833 }
2834
2835 #[test]
2836 fn scheduled_handlers_returns_only_scheduled() {
2837 let mut engine = create_test_engine();
2838 engine.register(EchoWorkflow).unwrap();
2839 engine.register(ScheduledWorkflow::new()).unwrap();
2840 engine.register(FailingWorkflow).unwrap();
2841
2842 let scheduled = engine.scheduled_handlers();
2843 assert_eq!(scheduled.len(), 1);
2844 assert_eq!(scheduled[0].0, "scheduled");
2845 assert_eq!(scheduled[0].1.as_str(), "0 0 * * * *");
2846 }
2847
2848 #[test]
2849 fn scheduled_handlers_empty_when_none_scheduled() {
2850 let mut engine = create_test_engine();
2851 engine.register(EchoWorkflow).unwrap();
2852 engine.register(FailingWorkflow).unwrap();
2853
2854 let scheduled = engine.scheduled_handlers();
2855 assert!(scheduled.is_empty());
2856 }
2857
2858 struct BadCategoryWorkflow(&'static str);
2859
2860 impl WorkflowHandler for BadCategoryWorkflow {
2861 fn name(&self) -> &str {
2862 "bad-category"
2863 }
2864 fn category(&self) -> Option<&str> {
2865 Some(self.0)
2866 }
2867 fn execute<'a>(
2868 &'a self,
2869 _ctx: &'a mut WorkflowContext,
2870 ) -> crate::handler::HandlerFuture<'a> {
2871 Box::pin(async move { Ok(()) })
2872 }
2873 }
2874
2875 #[test]
2876 fn engine_register_rejects_empty_category() {
2877 let mut engine = create_test_engine();
2878 let err = engine.register(BadCategoryWorkflow("")).unwrap_err();
2879 match err {
2880 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty category")),
2881 other => panic!("expected InvalidWorkflow, got {other:?}"),
2882 }
2883 }
2884
2885 #[test]
2886 fn engine_register_rejects_leading_slash_category() {
2887 let mut engine = create_test_engine();
2888 let err = engine
2889 .register(BadCategoryWorkflow("/data/etl"))
2890 .unwrap_err();
2891 match err {
2892 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("leading '/'")),
2893 other => panic!("expected InvalidWorkflow, got {other:?}"),
2894 }
2895 }
2896
2897 #[test]
2898 fn engine_register_rejects_trailing_slash_category() {
2899 let mut engine = create_test_engine();
2900 let err = engine
2901 .register(BadCategoryWorkflow("data/etl/"))
2902 .unwrap_err();
2903 match err {
2904 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("trailing '/'")),
2905 other => panic!("expected InvalidWorkflow, got {other:?}"),
2906 }
2907 }
2908
2909 #[test]
2910 fn engine_register_rejects_double_slash_category() {
2911 let mut engine = create_test_engine();
2912 let err = engine
2913 .register(BadCategoryWorkflow("data//etl"))
2914 .unwrap_err();
2915 match err {
2916 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty segment")),
2917 other => panic!("expected InvalidWorkflow, got {other:?}"),
2918 }
2919 }
2920
2921 #[test]
2922 fn engine_register_rejects_whitespace_only_segment_category() {
2923 let mut engine = create_test_engine();
2924 let err = engine
2925 .register(BadCategoryWorkflow("data/ /etl"))
2926 .unwrap_err();
2927 match err {
2928 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("whitespace-only segment")),
2929 other => panic!("expected InvalidWorkflow, got {other:?}"),
2930 }
2931 }
2932
2933 #[test]
2934 fn engine_register_accepts_valid_nested_category() {
2935 let mut engine = create_test_engine();
2936 assert!(engine.register(CategorizedWorkflow).is_ok());
2937 }
2938
2939 #[tokio::test]
2940 async fn engine_unknown_workflow_returns_error() {
2941 let engine = create_test_engine();
2942 let result = engine
2943 .run_handler("unknown", TriggerKind::Manual, json!({}))
2944 .await;
2945 assert!(result.is_err());
2946 match result {
2947 Err(EngineError::InvalidWorkflow(msg)) => {
2948 assert!(msg.contains("no handler registered"));
2949 }
2950 _ => panic!("expected InvalidWorkflow error"),
2951 }
2952 }
2953
2954 #[tokio::test]
2955 async fn engine_enqueue_handler_creates_pending_run() {
2956 let mut engine = create_test_engine();
2957 engine.register(EchoWorkflow).unwrap();
2958
2959 let run = engine
2960 .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
2961 .await
2962 .unwrap();
2963 assert_eq!(run.status.state, RunStatus::Pending);
2964 assert_eq!(run.workflow_name, "echo-workflow");
2965 }
2966
2967 #[tokio::test]
2968 async fn enqueue_handler_leaves_the_run_unattributed() {
2969 let mut engine = create_test_engine();
2970 engine.register(EchoWorkflow).unwrap();
2971
2972 let run = engine
2973 .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
2974 .await
2975 .unwrap();
2976
2977 assert!(run.created_by.is_none());
2978 }
2979
2980 #[tokio::test]
2981 async fn enqueue_handler_with_options_records_the_author() {
2982 let mut engine = create_test_engine();
2983 engine.register(EchoWorkflow).unwrap();
2984 let actor = RunActor::User {
2985 user_id: Uuid::now_v7(),
2986 };
2987
2988 let run = engine
2989 .enqueue_handler_with_options(
2990 "echo-workflow",
2991 TriggerKind::Api,
2992 json!({}),
2993 EnqueueOptions {
2994 created_by: Some(actor.clone()),
2995 ..Default::default()
2996 },
2997 )
2998 .await
2999 .unwrap()
3000 .into_run();
3001
3002 assert_eq!(run.created_by, Some(actor));
3003 }
3004
3005 #[tokio::test]
3006 async fn enqueue_handler_with_options_accepts_no_author() {
3007 let mut engine = create_test_engine();
3008 engine.register(EchoWorkflow).unwrap();
3009
3010 let run = engine
3011 .enqueue_handler_with_options(
3012 "echo-workflow",
3013 TriggerKind::Cron {
3014 schedule: "0 * * * * *".to_string(),
3015 },
3016 json!({}),
3017 EnqueueOptions::default(),
3018 )
3019 .await
3020 .unwrap()
3021 .into_run();
3022
3023 assert!(run.created_by.is_none());
3024 }
3025
3026 #[tokio::test]
3027 async fn enqueue_handler_with_options_stores_concurrency_limits() {
3028 let mut engine = create_test_engine();
3029 engine.register(EchoWorkflow).unwrap();
3030 let limits = vec![
3031 ConcurrencyLimit::new("repo:acme", 2),
3032 ConcurrencyLimit::new("tenant:42", 5),
3033 ];
3034
3035 let run = engine
3036 .enqueue_handler_with_options(
3037 "echo-workflow",
3038 TriggerKind::Api,
3039 json!({}),
3040 EnqueueOptions {
3041 concurrency_limits: limits.clone(),
3042 ..Default::default()
3043 },
3044 )
3045 .await
3046 .unwrap()
3047 .into_run();
3048
3049 assert_eq!(run.concurrency_limits, limits);
3050 }
3051
3052 #[tokio::test]
3053 async fn enqueue_rejects_invalid_concurrency_limits() {
3054 let mut engine = create_test_engine();
3055 engine.register(EchoWorkflow).unwrap();
3056
3057 let invalid = [
3058 vec![ConcurrencyLimit::new("repo:acme", 0)],
3059 vec![ConcurrencyLimit::new("", 1)],
3060 vec![
3061 ConcurrencyLimit::new("repo:acme", 1),
3062 ConcurrencyLimit::new("repo:acme", 2),
3063 ],
3064 ];
3065 for concurrency_limits in invalid {
3066 let err = engine
3067 .enqueue_handler_with_options(
3068 "echo-workflow",
3069 TriggerKind::Api,
3070 json!({}),
3071 EnqueueOptions {
3072 concurrency_limits,
3073 ..Default::default()
3074 },
3075 )
3076 .await
3077 .unwrap_err();
3078 assert!(
3079 matches!(err, EngineError::InvalidConcurrencyLimit(_)),
3080 "{err:?}"
3081 );
3082 }
3083
3084 let err = engine
3086 .enqueue_handler_with_options(
3087 "not-registered",
3088 TriggerKind::Api,
3089 json!({}),
3090 EnqueueOptions {
3091 concurrency_limits: vec![ConcurrencyLimit::new("repo:acme", 0)],
3092 ..Default::default()
3093 },
3094 )
3095 .await
3096 .unwrap_err();
3097 assert!(
3098 matches!(err, EngineError::InvalidConcurrencyLimit(_)),
3099 "{err:?}"
3100 );
3101
3102 let page = engine
3103 .store()
3104 .list_runs(RunFilter::default(), 1, 10)
3105 .await
3106 .unwrap();
3107 assert_eq!(page.total, 0, "no run may be created");
3108 }
3109
3110 struct UrgentWorkflow;
3111
3112 impl WorkflowHandler for UrgentWorkflow {
3113 fn name(&self) -> &str {
3114 "urgent-workflow"
3115 }
3116
3117 fn priority(&self) -> i16 {
3118 60
3119 }
3120
3121 fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3122 Box::pin(async { Ok(()) })
3123 }
3124 }
3125
3126 #[tokio::test]
3127 async fn enqueue_priority_defaults_to_the_handler_priority() {
3128 let mut engine = create_test_engine();
3129 engine.register(EchoWorkflow).unwrap();
3130 engine.register(UrgentWorkflow).unwrap();
3131
3132 let echo = engine
3133 .enqueue_handler_with_options(
3134 "echo-workflow",
3135 TriggerKind::Api,
3136 json!({}),
3137 EnqueueOptions::default(),
3138 )
3139 .await
3140 .unwrap()
3141 .into_run();
3142 assert_eq!(echo.priority, 0);
3143
3144 let urgent = engine
3145 .enqueue_handler_with_options(
3146 "urgent-workflow",
3147 TriggerKind::Api,
3148 json!({}),
3149 EnqueueOptions::default(),
3150 )
3151 .await
3152 .unwrap()
3153 .into_run();
3154 assert_eq!(urgent.priority, 60);
3155 }
3156
3157 #[tokio::test]
3158 async fn enqueue_priority_explicit_value_overrides_the_handler() {
3159 let mut engine = create_test_engine();
3160 engine.register(UrgentWorkflow).unwrap();
3161
3162 let run = engine
3163 .enqueue_handler_with_options(
3164 "urgent-workflow",
3165 TriggerKind::Api,
3166 json!({}),
3167 EnqueueOptions {
3168 priority: Some(-20),
3169 ..Default::default()
3170 },
3171 )
3172 .await
3173 .unwrap()
3174 .into_run();
3175 assert_eq!(run.priority, -20);
3176
3177 let stored = engine.store().get_run(run.id).await.unwrap().unwrap();
3178 assert_eq!(stored.priority, -20);
3179 }
3180
3181 #[tokio::test]
3182 async fn enqueue_priority_out_of_range_is_rejected() {
3183 let mut engine = create_test_engine();
3184 engine.register(EchoWorkflow).unwrap();
3185
3186 for priority in [MAX_PRIORITY + 1, MIN_PRIORITY - 1] {
3187 let err = engine
3188 .enqueue_handler_with_options(
3189 "echo-workflow",
3190 TriggerKind::Api,
3191 json!({}),
3192 EnqueueOptions {
3193 priority: Some(priority),
3194 ..Default::default()
3195 },
3196 )
3197 .await
3198 .unwrap_err();
3199 assert!(matches!(err, EngineError::InvalidPriority(_)), "{err:?}");
3200 }
3201
3202 let err = engine
3204 .enqueue_handler_with_options(
3205 "not-registered",
3206 TriggerKind::Api,
3207 json!({}),
3208 EnqueueOptions {
3209 priority: Some(MAX_PRIORITY + 1),
3210 ..Default::default()
3211 },
3212 )
3213 .await
3214 .unwrap_err();
3215 assert!(matches!(err, EngineError::InvalidPriority(_)), "{err:?}");
3216
3217 let page = engine
3218 .store()
3219 .list_runs(RunFilter::default(), 1, 10)
3220 .await
3221 .unwrap();
3222 assert_eq!(page.total, 0, "no run may be created");
3223 }
3224
3225 #[tokio::test]
3226 async fn enqueue_priority_bounds_are_accepted() {
3227 let mut engine = create_test_engine();
3228 engine.register(EchoWorkflow).unwrap();
3229
3230 for priority in [MIN_PRIORITY, MAX_PRIORITY] {
3231 let run = engine
3232 .enqueue_handler_with_options(
3233 "echo-workflow",
3234 TriggerKind::Api,
3235 json!({}),
3236 EnqueueOptions {
3237 priority: Some(priority),
3238 ..Default::default()
3239 },
3240 )
3241 .await
3242 .unwrap()
3243 .into_run();
3244 assert_eq!(run.priority, priority);
3245 }
3246 }
3247
3248 struct GpuWorkflow;
3249
3250 impl WorkflowHandler for GpuWorkflow {
3251 fn name(&self) -> &str {
3252 "gpu-workflow"
3253 }
3254
3255 fn required_worker_tags(&self) -> Vec<String> {
3256 vec!["gpu".to_string()]
3257 }
3258
3259 fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3260 Box::pin(async move { Ok(()) })
3261 }
3262 }
3263
3264 #[tokio::test]
3265 async fn enqueue_merges_handler_and_request_worker_tags() {
3266 let mut engine = create_test_engine();
3267 engine.register(GpuWorkflow).unwrap();
3268
3269 let run = engine
3270 .enqueue_handler_with_options(
3271 "gpu-workflow",
3272 TriggerKind::Api,
3273 json!({}),
3274 EnqueueOptions {
3275 worker_tags: vec!["region:eu".to_string(), "gpu".to_string()],
3276 ..Default::default()
3277 },
3278 )
3279 .await
3280 .unwrap()
3281 .into_run();
3282
3283 assert_eq!(
3284 run.worker_tags,
3285 vec!["gpu".to_string(), "region:eu".to_string()]
3286 );
3287 }
3288
3289 #[tokio::test]
3290 async fn enqueue_without_worker_tags_keeps_handler_tags() {
3291 let mut engine = create_test_engine();
3292 engine.register(GpuWorkflow).unwrap();
3293 engine.register(EchoWorkflow).unwrap();
3294
3295 let gpu = engine
3296 .enqueue_handler("gpu-workflow", TriggerKind::Api, json!({}), 0)
3297 .await
3298 .unwrap();
3299 assert_eq!(gpu.worker_tags, vec!["gpu".to_string()]);
3300
3301 let echo = engine
3302 .enqueue_handler("echo-workflow", TriggerKind::Api, json!({}), 0)
3303 .await
3304 .unwrap();
3305 assert!(echo.worker_tags.is_empty());
3306 }
3307
3308 #[tokio::test]
3309 async fn enqueue_rejects_invalid_worker_tags() {
3310 let mut engine = create_test_engine();
3311 engine.register(EchoWorkflow).unwrap();
3312
3313 for worker_tags in [
3314 vec!["bad,tag".to_string()],
3315 vec![" ".to_string()],
3316 vec!["x".repeat(65)],
3317 ] {
3318 let err = engine
3319 .enqueue_handler_with_options(
3320 "echo-workflow",
3321 TriggerKind::Api,
3322 json!({}),
3323 EnqueueOptions {
3324 worker_tags,
3325 ..Default::default()
3326 },
3327 )
3328 .await
3329 .unwrap_err();
3330 assert!(matches!(err, EngineError::InvalidWorkerTag(_)), "{err:?}");
3331 }
3332
3333 let err = engine
3335 .enqueue_handler_with_options(
3336 "not-registered",
3337 TriggerKind::Api,
3338 json!({}),
3339 EnqueueOptions {
3340 worker_tags: vec!["bad,tag".to_string()],
3341 ..Default::default()
3342 },
3343 )
3344 .await
3345 .unwrap_err();
3346 assert!(matches!(err, EngineError::InvalidWorkerTag(_)), "{err:?}");
3347 }
3348
3349 #[test]
3350 fn worker_tags_are_unset_by_default() {
3351 let engine = create_test_engine();
3352 assert!(engine.worker_tags().is_none());
3353 }
3354
3355 #[test]
3356 fn set_worker_tags_stores_the_tags() {
3357 let mut engine = create_test_engine();
3358 engine.set_worker_tags(vec!["gpu".to_string()]);
3359 assert_eq!(engine.worker_tags(), Some(&["gpu".to_string()][..]));
3360
3361 engine.set_worker_tags(Vec::new());
3362 assert_eq!(engine.worker_tags(), Some(&[][..]));
3363 }
3364
3365 #[tokio::test]
3366 async fn run_handler_records_handler_worker_tags() {
3367 let mut engine = create_test_engine();
3368 engine.register(GpuWorkflow).unwrap();
3369
3370 let result = engine
3371 .run_handler("gpu-workflow", TriggerKind::Manual, json!({}))
3372 .await
3373 .unwrap();
3374 assert_eq!(result.run.worker_tags, vec!["gpu".to_string()]);
3375 }
3376
3377 #[tokio::test]
3378 async fn run_handler_leaves_the_run_unattributed() {
3379 let mut engine = create_test_engine();
3380 engine.register(EchoWorkflow).unwrap();
3381
3382 let run = engine
3383 .run_handler("echo-workflow", TriggerKind::Manual, json!({}))
3384 .await
3385 .unwrap()
3386 .run;
3387
3388 assert!(run.created_by.is_none());
3389 }
3390
3391 #[tokio::test]
3392 async fn run_handler_priority_comes_from_the_handler() {
3393 let mut engine = create_test_engine();
3394 engine.register(UrgentWorkflow).unwrap();
3395
3396 let run = engine
3397 .run_handler("urgent-workflow", TriggerKind::Manual, json!({}))
3398 .await
3399 .unwrap()
3400 .run;
3401
3402 assert_eq!(run.priority, 60);
3403 }
3404
3405 #[tokio::test]
3406 async fn engine_register_boxed() {
3407 let mut engine = create_test_engine();
3408 let handler: Box<dyn WorkflowHandler> = Box::new(EchoWorkflow);
3409 let result = engine.register_boxed(handler);
3410 assert!(result.is_ok());
3411 assert_eq!(engine.handler_names().len(), 1);
3412 }
3413
3414 #[tokio::test]
3415 async fn engine_store_and_provider_accessors() {
3416 let store = Arc::new(InMemoryStore::new());
3417 let inner = ClaudeCodeProvider::new();
3418 let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
3419 inner,
3420 "/tmp/ironflow-fixtures",
3421 ));
3422 let engine = Engine::new(store.clone(), provider.clone());
3423
3424 let _ = engine.store();
3426 let _ = engine.provider();
3427 }
3428
3429 use crate::operation::{Operation, OperationContext};
3434 use async_trait::async_trait;
3435 use ironflow_core::error::OperationError;
3436 use ironflow_store::models::StepKind;
3437
3438 struct FakeGitlabOp {
3439 project_id: u64,
3440 title: String,
3441 }
3442
3443 #[async_trait]
3444 impl Operation for FakeGitlabOp {
3445 fn kind(&self) -> &str {
3446 "gitlab"
3447 }
3448
3449 async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
3450 Ok(json!({
3451 "issue_id": 42,
3452 "project_id": self.project_id,
3453 "title": self.title,
3454 }))
3455 }
3456
3457 fn input(&self) -> Option<Value> {
3458 Some(json!({
3459 "project_id": self.project_id,
3460 "title": self.title,
3461 }))
3462 }
3463 }
3464
3465 struct FailingOp;
3466
3467 #[async_trait]
3468 impl Operation for FailingOp {
3469 fn kind(&self) -> &str {
3470 "broken-service"
3471 }
3472
3473 async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
3474 Err(OperationError::Http {
3475 status: None,
3476 message: "service unavailable".to_string(),
3477 })
3478 }
3479 }
3480
3481 struct OperationWorkflow;
3482
3483 impl WorkflowHandler for OperationWorkflow {
3484 fn name(&self) -> &str {
3485 "operation-workflow"
3486 }
3487
3488 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3489 Box::pin(async move {
3490 let op = FakeGitlabOp {
3491 project_id: 123,
3492 title: "Bug report".to_string(),
3493 };
3494 ctx.operation("create-issue", &op).await?;
3495 Ok(())
3496 })
3497 }
3498 }
3499
3500 struct FailingOperationWorkflow;
3501
3502 impl WorkflowHandler for FailingOperationWorkflow {
3503 fn name(&self) -> &str {
3504 "failing-operation-workflow"
3505 }
3506
3507 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3508 Box::pin(async move {
3509 ctx.operation("broken-call", &FailingOp).await?;
3510 Ok(())
3511 })
3512 }
3513 }
3514
3515 struct MixedWorkflow;
3516
3517 impl WorkflowHandler for MixedWorkflow {
3518 fn name(&self) -> &str {
3519 "mixed-workflow"
3520 }
3521
3522 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3523 Box::pin(async move {
3524 ctx.shell("build", ShellConfig::new("echo built")).await?;
3525 let op = FakeGitlabOp {
3526 project_id: 456,
3527 title: "Deploy done".to_string(),
3528 };
3529 let result = ctx.operation("notify-gitlab", &op).await?;
3530 assert_eq!(result.output["issue_id"], 42);
3531 Ok(())
3532 })
3533 }
3534 }
3535
3536 #[tokio::test]
3537 async fn operation_step_happy_path() {
3538 let mut engine = create_test_engine();
3539 engine.register(OperationWorkflow).unwrap();
3540
3541 let run = engine
3542 .run_handler("operation-workflow", TriggerKind::Manual, json!({}))
3543 .await
3544 .unwrap()
3545 .run;
3546
3547 assert_eq!(run.status.state, RunStatus::Completed);
3548
3549 let steps = engine.store().list_steps(run.id).await.unwrap();
3550
3551 assert_eq!(steps.len(), 1);
3552 assert_eq!(steps[0].name, "create-issue");
3553 assert_eq!(steps[0].kind, StepKind::Custom("gitlab".to_string()));
3554 assert_eq!(
3555 steps[0].status.state,
3556 ironflow_store::models::StepStatus::Completed
3557 );
3558
3559 let output = steps[0].output.as_ref().unwrap();
3560 assert_eq!(output["issue_id"], 42);
3561 assert_eq!(output["project_id"], 123);
3562
3563 let input = steps[0].input.as_ref().unwrap();
3564 assert_eq!(input["project_id"], 123);
3565 assert_eq!(input["title"], "Bug report");
3566 }
3567
3568 #[tokio::test]
3569 async fn operation_step_failure_marks_run_failed() {
3570 let mut engine = create_test_engine();
3571 engine.register(FailingOperationWorkflow).unwrap();
3572
3573 let result = engine
3574 .run_handler("failing-operation-workflow", TriggerKind::Manual, json!({}))
3575 .await;
3576
3577 assert!(result.is_err());
3578 }
3579
3580 #[tokio::test]
3581 async fn operation_mixed_with_shell_steps() {
3582 let mut engine = create_test_engine();
3583 engine.register(MixedWorkflow).unwrap();
3584
3585 let run = engine
3586 .run_handler("mixed-workflow", TriggerKind::Manual, json!({}))
3587 .await
3588 .unwrap()
3589 .run;
3590
3591 assert_eq!(run.status.state, RunStatus::Completed);
3592
3593 let steps = engine.store().list_steps(run.id).await.unwrap();
3594
3595 assert_eq!(steps.len(), 2);
3596 assert_eq!(steps[0].kind, StepKind::Shell);
3597 assert_eq!(steps[1].kind, StepKind::Custom("gitlab".to_string()));
3598 assert_eq!(steps[0].position, 0);
3599 assert_eq!(steps[1].position, 1);
3600 }
3601
3602 use crate::config::ApprovalConfig;
3607
3608 struct SingleApprovalWorkflow;
3609
3610 impl WorkflowHandler for SingleApprovalWorkflow {
3611 fn name(&self) -> &str {
3612 "single-approval"
3613 }
3614
3615 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3616 Box::pin(async move {
3617 ctx.shell("build", ShellConfig::new("echo built")).await?;
3618 ctx.approval("gate", ApprovalConfig::new("OK?")).await?;
3619 ctx.shell("deploy", ShellConfig::new("echo deployed"))
3620 .await?;
3621 Ok(())
3622 })
3623 }
3624 }
3625
3626 struct DoubleApprovalWorkflow;
3627
3628 impl WorkflowHandler for DoubleApprovalWorkflow {
3629 fn name(&self) -> &str {
3630 "double-approval"
3631 }
3632
3633 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3634 Box::pin(async move {
3635 ctx.shell("build", ShellConfig::new("echo built")).await?;
3636 ctx.approval("staging-gate", ApprovalConfig::new("Deploy staging?"))
3637 .await?;
3638 ctx.shell("deploy-staging", ShellConfig::new("echo staging"))
3639 .await?;
3640 ctx.approval("prod-gate", ApprovalConfig::new("Deploy prod?"))
3641 .await?;
3642 ctx.shell("deploy-prod", ShellConfig::new("echo prod"))
3643 .await?;
3644 Ok(())
3645 })
3646 }
3647 }
3648
3649 #[tokio::test]
3650 async fn approval_pauses_run() {
3651 let mut engine = create_test_engine();
3652 engine.register(SingleApprovalWorkflow).unwrap();
3653
3654 let run = engine
3655 .run_handler("single-approval", TriggerKind::Manual, json!({}))
3656 .await
3657 .unwrap()
3658 .run;
3659
3660 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3661
3662 let steps = engine.store().list_steps(run.id).await.unwrap();
3663 assert_eq!(steps.len(), 2); assert_eq!(steps[0].kind, StepKind::Shell);
3665 assert_eq!(steps[0].status.state, StepStatus::Completed);
3666 assert_eq!(steps[1].kind, StepKind::Approval);
3667 assert_eq!(steps[1].status.state, StepStatus::AwaitingApproval);
3668 }
3669
3670 #[tokio::test]
3671 async fn approval_resume_completes_run() {
3672 let mut engine = create_test_engine();
3673 engine.register(SingleApprovalWorkflow).unwrap();
3674
3675 let run = engine
3677 .run_handler("single-approval", TriggerKind::Manual, json!({}))
3678 .await
3679 .unwrap()
3680 .run;
3681 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3682
3683 engine
3685 .store()
3686 .update_run_status(run.id, RunStatus::Running)
3687 .await
3688 .unwrap();
3689
3690 let resumed = engine.resume_run(run.id).await.unwrap().run;
3692 assert_eq!(resumed.status.state, RunStatus::Completed);
3693
3694 let steps = engine.store().list_steps(run.id).await.unwrap();
3695 assert_eq!(steps.len(), 3); assert_eq!(steps[0].name, "build");
3697 assert_eq!(steps[0].status.state, StepStatus::Completed);
3698 assert_eq!(steps[1].name, "gate");
3699 assert_eq!(steps[1].kind, StepKind::Approval);
3700 assert_eq!(steps[1].status.state, StepStatus::Completed);
3701 assert_eq!(steps[2].name, "deploy");
3702 assert_eq!(steps[2].status.state, StepStatus::Completed);
3703 }
3704
3705 #[tokio::test]
3706 async fn double_approval_two_resumes() {
3707 let mut engine = create_test_engine();
3708 engine.register(DoubleApprovalWorkflow).unwrap();
3709
3710 let run = engine
3712 .run_handler("double-approval", TriggerKind::Manual, json!({}))
3713 .await
3714 .unwrap()
3715 .run;
3716 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3717
3718 let steps = engine.store().list_steps(run.id).await.unwrap();
3719 assert_eq!(steps.len(), 2); engine
3723 .store()
3724 .update_run_status(run.id, RunStatus::Running)
3725 .await
3726 .unwrap();
3727
3728 let resumed = engine.resume_run(run.id).await.unwrap().run;
3729 assert_eq!(resumed.status.state, RunStatus::AwaitingApproval);
3730
3731 let steps = engine.store().list_steps(run.id).await.unwrap();
3732 assert_eq!(steps.len(), 4); engine
3736 .store()
3737 .update_run_status(run.id, RunStatus::Running)
3738 .await
3739 .unwrap();
3740
3741 let final_run = engine.resume_run(run.id).await.unwrap().run;
3742 assert_eq!(final_run.status.state, RunStatus::Completed);
3743
3744 let steps = engine.store().list_steps(run.id).await.unwrap();
3745 assert_eq!(steps.len(), 5);
3746 assert_eq!(steps[0].name, "build");
3747 assert_eq!(steps[1].name, "staging-gate");
3748 assert_eq!(steps[2].name, "deploy-staging");
3749 assert_eq!(steps[3].name, "prod-gate");
3750 assert_eq!(steps[4].name, "deploy-prod");
3751
3752 for step in &steps {
3753 assert_eq!(step.status.state, StepStatus::Completed);
3754 }
3755 }
3756
3757 use ironflow_store::models::{NewStep, StepUpdate, step_trace_id};
3762
3763 async fn create_step_with_status(
3764 store: &Arc<dyn Store>,
3765 run_id: Uuid,
3766 name: &str,
3767 position: u32,
3768 status: StepStatus,
3769 ) -> ironflow_store::models::Step {
3770 let step = store
3771 .create_step(NewStep {
3772 run_id,
3773 trace_id: step_trace_id(run_id, name, position),
3774 name: name.to_string(),
3775 kind: StepKind::Shell,
3776 position,
3777 input: None,
3778 is_error_handler: false,
3779 })
3780 .await
3781 .unwrap();
3782
3783 match status {
3784 StepStatus::Pending => {}
3785 StepStatus::Running => {
3786 store
3787 .update_step(
3788 step.id,
3789 StepUpdate {
3790 status: Some(StepStatus::Running),
3791 ..StepUpdate::default()
3792 },
3793 )
3794 .await
3795 .unwrap();
3796 }
3797 StepStatus::Completed => {
3798 store
3799 .update_step(
3800 step.id,
3801 StepUpdate {
3802 status: Some(StepStatus::Running),
3803 ..StepUpdate::default()
3804 },
3805 )
3806 .await
3807 .unwrap();
3808 store
3809 .update_step(
3810 step.id,
3811 StepUpdate {
3812 status: Some(StepStatus::Completed),
3813 ..StepUpdate::default()
3814 },
3815 )
3816 .await
3817 .unwrap();
3818 }
3819 StepStatus::AwaitingApproval => {
3820 store
3821 .update_step(
3822 step.id,
3823 StepUpdate {
3824 status: Some(StepStatus::Running),
3825 ..StepUpdate::default()
3826 },
3827 )
3828 .await
3829 .unwrap();
3830 store
3831 .update_step(
3832 step.id,
3833 StepUpdate {
3834 status: Some(StepStatus::AwaitingApproval),
3835 ..StepUpdate::default()
3836 },
3837 )
3838 .await
3839 .unwrap();
3840 }
3841 _ => panic!("unsupported status for test helper: {status}"),
3842 }
3843
3844 store.get_step(step.id).await.unwrap().unwrap()
3845 }
3846
3847 #[tokio::test]
3848 async fn fail_orphaned_steps_marks_running_as_failed() {
3849 let engine = create_test_engine();
3850 let run = engine
3851 .store()
3852 .create_run(NewRun {
3853 created_by: None,
3854 workflow_name: "test".to_string(),
3855 trigger: TriggerKind::Manual,
3856 payload: json!({}),
3857 max_retries: 0,
3858 handler_version: None,
3859 labels: HashMap::new(),
3860 scheduled_at: None,
3861 idempotency_key: None,
3862 concurrency_key: None,
3863 priority: 0,
3864 concurrency_limits: Vec::new(),
3865 max_cost_usd: None,
3866 worker_tags: Vec::new(),
3867 })
3868 .await
3869 .unwrap()
3870 .into_run();
3871
3872 let step = create_step_with_status(
3873 engine.store(),
3874 run.id,
3875 "running-step",
3876 0,
3877 StepStatus::Running,
3878 )
3879 .await;
3880
3881 engine
3882 .fail_orphaned_steps(run.id, "parent run timed out")
3883 .await
3884 .unwrap();
3885
3886 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
3887 assert_eq!(updated.status.state, StepStatus::Failed);
3888 assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
3889 assert!(updated.completed_at.is_some());
3890 }
3891
3892 #[tokio::test]
3893 async fn fail_orphaned_steps_marks_pending_as_skipped() {
3894 let engine = create_test_engine();
3895 let run = engine
3896 .store()
3897 .create_run(NewRun {
3898 created_by: None,
3899 workflow_name: "test".to_string(),
3900 trigger: TriggerKind::Manual,
3901 payload: json!({}),
3902 max_retries: 0,
3903 handler_version: None,
3904 labels: HashMap::new(),
3905 scheduled_at: None,
3906 idempotency_key: None,
3907 concurrency_key: None,
3908 priority: 0,
3909 concurrency_limits: Vec::new(),
3910 max_cost_usd: None,
3911 worker_tags: Vec::new(),
3912 })
3913 .await
3914 .unwrap()
3915 .into_run();
3916
3917 let step = create_step_with_status(
3918 engine.store(),
3919 run.id,
3920 "pending-step",
3921 0,
3922 StepStatus::Pending,
3923 )
3924 .await;
3925
3926 engine
3927 .fail_orphaned_steps(run.id, "parent run timed out")
3928 .await
3929 .unwrap();
3930
3931 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
3932 assert_eq!(updated.status.state, StepStatus::Skipped);
3933 assert!(updated.error.is_none());
3934 assert!(updated.completed_at.is_some());
3935 }
3936
3937 #[tokio::test]
3938 async fn fail_orphaned_steps_marks_awaiting_approval_as_failed() {
3939 let engine = create_test_engine();
3940 let run = engine
3941 .store()
3942 .create_run(NewRun {
3943 created_by: None,
3944 workflow_name: "test".to_string(),
3945 trigger: TriggerKind::Manual,
3946 payload: json!({}),
3947 max_retries: 0,
3948 handler_version: None,
3949 labels: HashMap::new(),
3950 scheduled_at: None,
3951 idempotency_key: None,
3952 concurrency_key: None,
3953 priority: 0,
3954 concurrency_limits: Vec::new(),
3955 max_cost_usd: None,
3956 worker_tags: Vec::new(),
3957 })
3958 .await
3959 .unwrap()
3960 .into_run();
3961
3962 let step = create_step_with_status(
3963 engine.store(),
3964 run.id,
3965 "approval-step",
3966 0,
3967 StepStatus::AwaitingApproval,
3968 )
3969 .await;
3970
3971 engine
3972 .fail_orphaned_steps(run.id, "parent run timed out")
3973 .await
3974 .unwrap();
3975
3976 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
3977 assert_eq!(updated.status.state, StepStatus::Failed);
3978 assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
3979 assert!(updated.completed_at.is_some());
3980 }
3981
3982 #[tokio::test]
3983 async fn fail_orphaned_steps_skips_terminal_steps() {
3984 let engine = create_test_engine();
3985 let run = engine
3986 .store()
3987 .create_run(NewRun {
3988 created_by: None,
3989 workflow_name: "test".to_string(),
3990 trigger: TriggerKind::Manual,
3991 payload: json!({}),
3992 max_retries: 0,
3993 handler_version: None,
3994 labels: HashMap::new(),
3995 scheduled_at: None,
3996 idempotency_key: None,
3997 concurrency_key: None,
3998 priority: 0,
3999 concurrency_limits: Vec::new(),
4000 max_cost_usd: None,
4001 worker_tags: Vec::new(),
4002 })
4003 .await
4004 .unwrap()
4005 .into_run();
4006
4007 let completed_step =
4008 create_step_with_status(engine.store(), run.id, "done", 0, StepStatus::Completed).await;
4009 let running_step =
4010 create_step_with_status(engine.store(), run.id, "in-flight", 1, StepStatus::Running)
4011 .await;
4012
4013 engine
4014 .fail_orphaned_steps(run.id, "parent run timed out")
4015 .await
4016 .unwrap();
4017
4018 let completed = engine
4019 .store()
4020 .get_step(completed_step.id)
4021 .await
4022 .unwrap()
4023 .unwrap();
4024 assert_eq!(completed.status.state, StepStatus::Completed);
4025
4026 let failed = engine
4027 .store()
4028 .get_step(running_step.id)
4029 .await
4030 .unwrap()
4031 .unwrap();
4032 assert_eq!(failed.status.state, StepStatus::Failed);
4033 }
4034
4035 #[tokio::test]
4036 async fn fail_orphaned_steps_mixed_states() {
4037 let engine = create_test_engine();
4038 let run = engine
4039 .store()
4040 .create_run(NewRun {
4041 created_by: None,
4042 workflow_name: "test".to_string(),
4043 trigger: TriggerKind::Manual,
4044 payload: json!({}),
4045 max_retries: 0,
4046 handler_version: None,
4047 labels: HashMap::new(),
4048 scheduled_at: None,
4049 idempotency_key: None,
4050 concurrency_key: None,
4051 priority: 0,
4052 concurrency_limits: Vec::new(),
4053 max_cost_usd: None,
4054 worker_tags: Vec::new(),
4055 })
4056 .await
4057 .unwrap()
4058 .into_run();
4059
4060 let s_completed =
4061 create_step_with_status(engine.store(), run.id, "step-1", 0, StepStatus::Completed)
4062 .await;
4063 let s_running =
4064 create_step_with_status(engine.store(), run.id, "step-2", 1, StepStatus::Running).await;
4065 let s_pending =
4066 create_step_with_status(engine.store(), run.id, "step-3", 2, StepStatus::Pending).await;
4067
4068 engine.fail_orphaned_steps(run.id, "timeout").await.unwrap();
4069
4070 let r_completed = engine
4071 .store()
4072 .get_step(s_completed.id)
4073 .await
4074 .unwrap()
4075 .unwrap();
4076 assert_eq!(r_completed.status.state, StepStatus::Completed);
4077
4078 let r_running = engine
4079 .store()
4080 .get_step(s_running.id)
4081 .await
4082 .unwrap()
4083 .unwrap();
4084 assert_eq!(r_running.status.state, StepStatus::Failed);
4085 assert_eq!(r_running.error.as_deref(), Some("timeout"));
4086
4087 let r_pending = engine
4088 .store()
4089 .get_step(s_pending.id)
4090 .await
4091 .unwrap()
4092 .unwrap();
4093 assert_eq!(r_pending.status.state, StepStatus::Skipped);
4094 assert!(r_pending.error.is_none());
4095 }
4096
4097 #[tokio::test]
4098 async fn fail_orphaned_steps_no_steps_is_noop() {
4099 let engine = create_test_engine();
4100 let run = engine
4101 .store()
4102 .create_run(NewRun {
4103 created_by: None,
4104 workflow_name: "test".to_string(),
4105 trigger: TriggerKind::Manual,
4106 payload: json!({}),
4107 max_retries: 0,
4108 handler_version: None,
4109 labels: HashMap::new(),
4110 scheduled_at: None,
4111 idempotency_key: None,
4112 concurrency_key: None,
4113 priority: 0,
4114 concurrency_limits: Vec::new(),
4115 max_cost_usd: None,
4116 worker_tags: Vec::new(),
4117 })
4118 .await
4119 .unwrap()
4120 .into_run();
4121
4122 let result = engine.fail_orphaned_steps(run.id, "timeout").await;
4123 assert!(result.is_ok());
4124 }
4125
4126 #[tokio::test]
4127 async fn fail_orphaned_steps_preserves_existing_error() {
4128 let engine = create_test_engine();
4129 let run = engine
4130 .store()
4131 .create_run(NewRun {
4132 created_by: None,
4133 workflow_name: "test".to_string(),
4134 trigger: TriggerKind::Manual,
4135 payload: json!({}),
4136 max_retries: 0,
4137 handler_version: None,
4138 labels: HashMap::new(),
4139 scheduled_at: None,
4140 idempotency_key: None,
4141 concurrency_key: None,
4142 priority: 0,
4143 concurrency_limits: Vec::new(),
4144 max_cost_usd: None,
4145 worker_tags: Vec::new(),
4146 })
4147 .await
4148 .unwrap()
4149 .into_run();
4150
4151 let step_with_error = create_step_with_status(
4152 engine.store(),
4153 run.id,
4154 "already-errored",
4155 0,
4156 StepStatus::Running,
4157 )
4158 .await;
4159
4160 engine
4161 .store()
4162 .update_step(
4163 step_with_error.id,
4164 StepUpdate {
4165 error: Some("real error from provider".to_string()),
4166 ..StepUpdate::default()
4167 },
4168 )
4169 .await
4170 .unwrap();
4171
4172 let step_no_error = create_step_with_status(
4173 engine.store(),
4174 run.id,
4175 "no-error-yet",
4176 1,
4177 StepStatus::Running,
4178 )
4179 .await;
4180
4181 engine
4182 .fail_orphaned_steps(run.id, "parent run failed")
4183 .await
4184 .unwrap();
4185
4186 let updated_with = engine
4187 .store()
4188 .get_step(step_with_error.id)
4189 .await
4190 .unwrap()
4191 .unwrap();
4192 assert_eq!(updated_with.status.state, StepStatus::Failed);
4193 assert_eq!(
4194 updated_with.error.as_deref(),
4195 Some("real error from provider"),
4196 );
4197
4198 let updated_without = engine
4199 .store()
4200 .get_step(step_no_error.id)
4201 .await
4202 .unwrap()
4203 .unwrap();
4204 assert_eq!(updated_without.status.state, StepStatus::Failed);
4205 assert_eq!(updated_without.error.as_deref(), Some("parent run failed"),);
4206 }
4207}