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 schedule_id: None,
3016 scheduled_for: None,
3017 },
3018 json!({}),
3019 EnqueueOptions::default(),
3020 )
3021 .await
3022 .unwrap()
3023 .into_run();
3024
3025 assert!(run.created_by.is_none());
3026 }
3027
3028 #[tokio::test]
3029 async fn enqueue_handler_with_options_stores_concurrency_limits() {
3030 let mut engine = create_test_engine();
3031 engine.register(EchoWorkflow).unwrap();
3032 let limits = vec![
3033 ConcurrencyLimit::new("repo:acme", 2),
3034 ConcurrencyLimit::new("tenant:42", 5),
3035 ];
3036
3037 let run = engine
3038 .enqueue_handler_with_options(
3039 "echo-workflow",
3040 TriggerKind::Api,
3041 json!({}),
3042 EnqueueOptions {
3043 concurrency_limits: limits.clone(),
3044 ..Default::default()
3045 },
3046 )
3047 .await
3048 .unwrap()
3049 .into_run();
3050
3051 assert_eq!(run.concurrency_limits, limits);
3052 }
3053
3054 #[tokio::test]
3055 async fn enqueue_rejects_invalid_concurrency_limits() {
3056 let mut engine = create_test_engine();
3057 engine.register(EchoWorkflow).unwrap();
3058
3059 let invalid = [
3060 vec![ConcurrencyLimit::new("repo:acme", 0)],
3061 vec![ConcurrencyLimit::new("", 1)],
3062 vec![
3063 ConcurrencyLimit::new("repo:acme", 1),
3064 ConcurrencyLimit::new("repo:acme", 2),
3065 ],
3066 ];
3067 for concurrency_limits in invalid {
3068 let err = engine
3069 .enqueue_handler_with_options(
3070 "echo-workflow",
3071 TriggerKind::Api,
3072 json!({}),
3073 EnqueueOptions {
3074 concurrency_limits,
3075 ..Default::default()
3076 },
3077 )
3078 .await
3079 .unwrap_err();
3080 assert!(
3081 matches!(err, EngineError::InvalidConcurrencyLimit(_)),
3082 "{err:?}"
3083 );
3084 }
3085
3086 let err = engine
3088 .enqueue_handler_with_options(
3089 "not-registered",
3090 TriggerKind::Api,
3091 json!({}),
3092 EnqueueOptions {
3093 concurrency_limits: vec![ConcurrencyLimit::new("repo:acme", 0)],
3094 ..Default::default()
3095 },
3096 )
3097 .await
3098 .unwrap_err();
3099 assert!(
3100 matches!(err, EngineError::InvalidConcurrencyLimit(_)),
3101 "{err:?}"
3102 );
3103
3104 let page = engine
3105 .store()
3106 .list_runs(RunFilter::default(), 1, 10)
3107 .await
3108 .unwrap();
3109 assert_eq!(page.total, 0, "no run may be created");
3110 }
3111
3112 struct UrgentWorkflow;
3113
3114 impl WorkflowHandler for UrgentWorkflow {
3115 fn name(&self) -> &str {
3116 "urgent-workflow"
3117 }
3118
3119 fn priority(&self) -> i16 {
3120 60
3121 }
3122
3123 fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3124 Box::pin(async { Ok(()) })
3125 }
3126 }
3127
3128 #[tokio::test]
3129 async fn enqueue_priority_defaults_to_the_handler_priority() {
3130 let mut engine = create_test_engine();
3131 engine.register(EchoWorkflow).unwrap();
3132 engine.register(UrgentWorkflow).unwrap();
3133
3134 let echo = engine
3135 .enqueue_handler_with_options(
3136 "echo-workflow",
3137 TriggerKind::Api,
3138 json!({}),
3139 EnqueueOptions::default(),
3140 )
3141 .await
3142 .unwrap()
3143 .into_run();
3144 assert_eq!(echo.priority, 0);
3145
3146 let urgent = engine
3147 .enqueue_handler_with_options(
3148 "urgent-workflow",
3149 TriggerKind::Api,
3150 json!({}),
3151 EnqueueOptions::default(),
3152 )
3153 .await
3154 .unwrap()
3155 .into_run();
3156 assert_eq!(urgent.priority, 60);
3157 }
3158
3159 #[tokio::test]
3160 async fn enqueue_priority_explicit_value_overrides_the_handler() {
3161 let mut engine = create_test_engine();
3162 engine.register(UrgentWorkflow).unwrap();
3163
3164 let run = engine
3165 .enqueue_handler_with_options(
3166 "urgent-workflow",
3167 TriggerKind::Api,
3168 json!({}),
3169 EnqueueOptions {
3170 priority: Some(-20),
3171 ..Default::default()
3172 },
3173 )
3174 .await
3175 .unwrap()
3176 .into_run();
3177 assert_eq!(run.priority, -20);
3178
3179 let stored = engine.store().get_run(run.id).await.unwrap().unwrap();
3180 assert_eq!(stored.priority, -20);
3181 }
3182
3183 #[tokio::test]
3184 async fn enqueue_priority_out_of_range_is_rejected() {
3185 let mut engine = create_test_engine();
3186 engine.register(EchoWorkflow).unwrap();
3187
3188 for priority in [MAX_PRIORITY + 1, MIN_PRIORITY - 1] {
3189 let err = engine
3190 .enqueue_handler_with_options(
3191 "echo-workflow",
3192 TriggerKind::Api,
3193 json!({}),
3194 EnqueueOptions {
3195 priority: Some(priority),
3196 ..Default::default()
3197 },
3198 )
3199 .await
3200 .unwrap_err();
3201 assert!(matches!(err, EngineError::InvalidPriority(_)), "{err:?}");
3202 }
3203
3204 let err = engine
3206 .enqueue_handler_with_options(
3207 "not-registered",
3208 TriggerKind::Api,
3209 json!({}),
3210 EnqueueOptions {
3211 priority: Some(MAX_PRIORITY + 1),
3212 ..Default::default()
3213 },
3214 )
3215 .await
3216 .unwrap_err();
3217 assert!(matches!(err, EngineError::InvalidPriority(_)), "{err:?}");
3218
3219 let page = engine
3220 .store()
3221 .list_runs(RunFilter::default(), 1, 10)
3222 .await
3223 .unwrap();
3224 assert_eq!(page.total, 0, "no run may be created");
3225 }
3226
3227 #[tokio::test]
3228 async fn enqueue_priority_bounds_are_accepted() {
3229 let mut engine = create_test_engine();
3230 engine.register(EchoWorkflow).unwrap();
3231
3232 for priority in [MIN_PRIORITY, MAX_PRIORITY] {
3233 let run = engine
3234 .enqueue_handler_with_options(
3235 "echo-workflow",
3236 TriggerKind::Api,
3237 json!({}),
3238 EnqueueOptions {
3239 priority: Some(priority),
3240 ..Default::default()
3241 },
3242 )
3243 .await
3244 .unwrap()
3245 .into_run();
3246 assert_eq!(run.priority, priority);
3247 }
3248 }
3249
3250 struct GpuWorkflow;
3251
3252 impl WorkflowHandler for GpuWorkflow {
3253 fn name(&self) -> &str {
3254 "gpu-workflow"
3255 }
3256
3257 fn required_worker_tags(&self) -> Vec<String> {
3258 vec!["gpu".to_string()]
3259 }
3260
3261 fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3262 Box::pin(async move { Ok(()) })
3263 }
3264 }
3265
3266 #[tokio::test]
3267 async fn enqueue_merges_handler_and_request_worker_tags() {
3268 let mut engine = create_test_engine();
3269 engine.register(GpuWorkflow).unwrap();
3270
3271 let run = engine
3272 .enqueue_handler_with_options(
3273 "gpu-workflow",
3274 TriggerKind::Api,
3275 json!({}),
3276 EnqueueOptions {
3277 worker_tags: vec!["region:eu".to_string(), "gpu".to_string()],
3278 ..Default::default()
3279 },
3280 )
3281 .await
3282 .unwrap()
3283 .into_run();
3284
3285 assert_eq!(
3286 run.worker_tags,
3287 vec!["gpu".to_string(), "region:eu".to_string()]
3288 );
3289 }
3290
3291 #[tokio::test]
3292 async fn enqueue_without_worker_tags_keeps_handler_tags() {
3293 let mut engine = create_test_engine();
3294 engine.register(GpuWorkflow).unwrap();
3295 engine.register(EchoWorkflow).unwrap();
3296
3297 let gpu = engine
3298 .enqueue_handler("gpu-workflow", TriggerKind::Api, json!({}), 0)
3299 .await
3300 .unwrap();
3301 assert_eq!(gpu.worker_tags, vec!["gpu".to_string()]);
3302
3303 let echo = engine
3304 .enqueue_handler("echo-workflow", TriggerKind::Api, json!({}), 0)
3305 .await
3306 .unwrap();
3307 assert!(echo.worker_tags.is_empty());
3308 }
3309
3310 #[tokio::test]
3311 async fn enqueue_rejects_invalid_worker_tags() {
3312 let mut engine = create_test_engine();
3313 engine.register(EchoWorkflow).unwrap();
3314
3315 for worker_tags in [
3316 vec!["bad,tag".to_string()],
3317 vec![" ".to_string()],
3318 vec!["x".repeat(65)],
3319 ] {
3320 let err = engine
3321 .enqueue_handler_with_options(
3322 "echo-workflow",
3323 TriggerKind::Api,
3324 json!({}),
3325 EnqueueOptions {
3326 worker_tags,
3327 ..Default::default()
3328 },
3329 )
3330 .await
3331 .unwrap_err();
3332 assert!(matches!(err, EngineError::InvalidWorkerTag(_)), "{err:?}");
3333 }
3334
3335 let err = engine
3337 .enqueue_handler_with_options(
3338 "not-registered",
3339 TriggerKind::Api,
3340 json!({}),
3341 EnqueueOptions {
3342 worker_tags: vec!["bad,tag".to_string()],
3343 ..Default::default()
3344 },
3345 )
3346 .await
3347 .unwrap_err();
3348 assert!(matches!(err, EngineError::InvalidWorkerTag(_)), "{err:?}");
3349 }
3350
3351 #[test]
3352 fn worker_tags_are_unset_by_default() {
3353 let engine = create_test_engine();
3354 assert!(engine.worker_tags().is_none());
3355 }
3356
3357 #[test]
3358 fn set_worker_tags_stores_the_tags() {
3359 let mut engine = create_test_engine();
3360 engine.set_worker_tags(vec!["gpu".to_string()]);
3361 assert_eq!(engine.worker_tags(), Some(&["gpu".to_string()][..]));
3362
3363 engine.set_worker_tags(Vec::new());
3364 assert_eq!(engine.worker_tags(), Some(&[][..]));
3365 }
3366
3367 #[tokio::test]
3368 async fn run_handler_records_handler_worker_tags() {
3369 let mut engine = create_test_engine();
3370 engine.register(GpuWorkflow).unwrap();
3371
3372 let result = engine
3373 .run_handler("gpu-workflow", TriggerKind::Manual, json!({}))
3374 .await
3375 .unwrap();
3376 assert_eq!(result.run.worker_tags, vec!["gpu".to_string()]);
3377 }
3378
3379 #[tokio::test]
3380 async fn run_handler_leaves_the_run_unattributed() {
3381 let mut engine = create_test_engine();
3382 engine.register(EchoWorkflow).unwrap();
3383
3384 let run = engine
3385 .run_handler("echo-workflow", TriggerKind::Manual, json!({}))
3386 .await
3387 .unwrap()
3388 .run;
3389
3390 assert!(run.created_by.is_none());
3391 }
3392
3393 #[tokio::test]
3394 async fn run_handler_priority_comes_from_the_handler() {
3395 let mut engine = create_test_engine();
3396 engine.register(UrgentWorkflow).unwrap();
3397
3398 let run = engine
3399 .run_handler("urgent-workflow", TriggerKind::Manual, json!({}))
3400 .await
3401 .unwrap()
3402 .run;
3403
3404 assert_eq!(run.priority, 60);
3405 }
3406
3407 #[tokio::test]
3408 async fn engine_register_boxed() {
3409 let mut engine = create_test_engine();
3410 let handler: Box<dyn WorkflowHandler> = Box::new(EchoWorkflow);
3411 let result = engine.register_boxed(handler);
3412 assert!(result.is_ok());
3413 assert_eq!(engine.handler_names().len(), 1);
3414 }
3415
3416 #[tokio::test]
3417 async fn engine_store_and_provider_accessors() {
3418 let store = Arc::new(InMemoryStore::new());
3419 let inner = ClaudeCodeProvider::new();
3420 let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
3421 inner,
3422 "/tmp/ironflow-fixtures",
3423 ));
3424 let engine = Engine::new(store.clone(), provider.clone());
3425
3426 let _ = engine.store();
3428 let _ = engine.provider();
3429 }
3430
3431 use crate::operation::{Operation, OperationContext};
3436 use async_trait::async_trait;
3437 use ironflow_core::error::OperationError;
3438 use ironflow_store::models::StepKind;
3439
3440 struct FakeGitlabOp {
3441 project_id: u64,
3442 title: String,
3443 }
3444
3445 #[async_trait]
3446 impl Operation for FakeGitlabOp {
3447 fn kind(&self) -> &str {
3448 "gitlab"
3449 }
3450
3451 async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
3452 Ok(json!({
3453 "issue_id": 42,
3454 "project_id": self.project_id,
3455 "title": self.title,
3456 }))
3457 }
3458
3459 fn input(&self) -> Option<Value> {
3460 Some(json!({
3461 "project_id": self.project_id,
3462 "title": self.title,
3463 }))
3464 }
3465 }
3466
3467 struct FailingOp;
3468
3469 #[async_trait]
3470 impl Operation for FailingOp {
3471 fn kind(&self) -> &str {
3472 "broken-service"
3473 }
3474
3475 async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
3476 Err(OperationError::Http {
3477 status: None,
3478 message: "service unavailable".to_string(),
3479 })
3480 }
3481 }
3482
3483 struct OperationWorkflow;
3484
3485 impl WorkflowHandler for OperationWorkflow {
3486 fn name(&self) -> &str {
3487 "operation-workflow"
3488 }
3489
3490 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3491 Box::pin(async move {
3492 let op = FakeGitlabOp {
3493 project_id: 123,
3494 title: "Bug report".to_string(),
3495 };
3496 ctx.operation("create-issue", &op).await?;
3497 Ok(())
3498 })
3499 }
3500 }
3501
3502 struct FailingOperationWorkflow;
3503
3504 impl WorkflowHandler for FailingOperationWorkflow {
3505 fn name(&self) -> &str {
3506 "failing-operation-workflow"
3507 }
3508
3509 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3510 Box::pin(async move {
3511 ctx.operation("broken-call", &FailingOp).await?;
3512 Ok(())
3513 })
3514 }
3515 }
3516
3517 struct MixedWorkflow;
3518
3519 impl WorkflowHandler for MixedWorkflow {
3520 fn name(&self) -> &str {
3521 "mixed-workflow"
3522 }
3523
3524 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3525 Box::pin(async move {
3526 ctx.shell("build", ShellConfig::new("echo built")).await?;
3527 let op = FakeGitlabOp {
3528 project_id: 456,
3529 title: "Deploy done".to_string(),
3530 };
3531 let result = ctx.operation("notify-gitlab", &op).await?;
3532 assert_eq!(result.output["issue_id"], 42);
3533 Ok(())
3534 })
3535 }
3536 }
3537
3538 #[tokio::test]
3539 async fn operation_step_happy_path() {
3540 let mut engine = create_test_engine();
3541 engine.register(OperationWorkflow).unwrap();
3542
3543 let run = engine
3544 .run_handler("operation-workflow", TriggerKind::Manual, json!({}))
3545 .await
3546 .unwrap()
3547 .run;
3548
3549 assert_eq!(run.status.state, RunStatus::Completed);
3550
3551 let steps = engine.store().list_steps(run.id).await.unwrap();
3552
3553 assert_eq!(steps.len(), 1);
3554 assert_eq!(steps[0].name, "create-issue");
3555 assert_eq!(steps[0].kind, StepKind::Custom("gitlab".to_string()));
3556 assert_eq!(
3557 steps[0].status.state,
3558 ironflow_store::models::StepStatus::Completed
3559 );
3560
3561 let output = steps[0].output.as_ref().unwrap();
3562 assert_eq!(output["issue_id"], 42);
3563 assert_eq!(output["project_id"], 123);
3564
3565 let input = steps[0].input.as_ref().unwrap();
3566 assert_eq!(input["project_id"], 123);
3567 assert_eq!(input["title"], "Bug report");
3568 }
3569
3570 #[tokio::test]
3571 async fn operation_step_failure_marks_run_failed() {
3572 let mut engine = create_test_engine();
3573 engine.register(FailingOperationWorkflow).unwrap();
3574
3575 let result = engine
3576 .run_handler("failing-operation-workflow", TriggerKind::Manual, json!({}))
3577 .await;
3578
3579 assert!(result.is_err());
3580 }
3581
3582 #[tokio::test]
3583 async fn operation_mixed_with_shell_steps() {
3584 let mut engine = create_test_engine();
3585 engine.register(MixedWorkflow).unwrap();
3586
3587 let run = engine
3588 .run_handler("mixed-workflow", TriggerKind::Manual, json!({}))
3589 .await
3590 .unwrap()
3591 .run;
3592
3593 assert_eq!(run.status.state, RunStatus::Completed);
3594
3595 let steps = engine.store().list_steps(run.id).await.unwrap();
3596
3597 assert_eq!(steps.len(), 2);
3598 assert_eq!(steps[0].kind, StepKind::Shell);
3599 assert_eq!(steps[1].kind, StepKind::Custom("gitlab".to_string()));
3600 assert_eq!(steps[0].position, 0);
3601 assert_eq!(steps[1].position, 1);
3602 }
3603
3604 use crate::config::ApprovalConfig;
3609
3610 struct SingleApprovalWorkflow;
3611
3612 impl WorkflowHandler for SingleApprovalWorkflow {
3613 fn name(&self) -> &str {
3614 "single-approval"
3615 }
3616
3617 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3618 Box::pin(async move {
3619 ctx.shell("build", ShellConfig::new("echo built")).await?;
3620 ctx.approval("gate", ApprovalConfig::new("OK?")).await?;
3621 ctx.shell("deploy", ShellConfig::new("echo deployed"))
3622 .await?;
3623 Ok(())
3624 })
3625 }
3626 }
3627
3628 struct DoubleApprovalWorkflow;
3629
3630 impl WorkflowHandler for DoubleApprovalWorkflow {
3631 fn name(&self) -> &str {
3632 "double-approval"
3633 }
3634
3635 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3636 Box::pin(async move {
3637 ctx.shell("build", ShellConfig::new("echo built")).await?;
3638 ctx.approval("staging-gate", ApprovalConfig::new("Deploy staging?"))
3639 .await?;
3640 ctx.shell("deploy-staging", ShellConfig::new("echo staging"))
3641 .await?;
3642 ctx.approval("prod-gate", ApprovalConfig::new("Deploy prod?"))
3643 .await?;
3644 ctx.shell("deploy-prod", ShellConfig::new("echo prod"))
3645 .await?;
3646 Ok(())
3647 })
3648 }
3649 }
3650
3651 #[tokio::test]
3652 async fn approval_pauses_run() {
3653 let mut engine = create_test_engine();
3654 engine.register(SingleApprovalWorkflow).unwrap();
3655
3656 let run = engine
3657 .run_handler("single-approval", TriggerKind::Manual, json!({}))
3658 .await
3659 .unwrap()
3660 .run;
3661
3662 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3663
3664 let steps = engine.store().list_steps(run.id).await.unwrap();
3665 assert_eq!(steps.len(), 2); assert_eq!(steps[0].kind, StepKind::Shell);
3667 assert_eq!(steps[0].status.state, StepStatus::Completed);
3668 assert_eq!(steps[1].kind, StepKind::Approval);
3669 assert_eq!(steps[1].status.state, StepStatus::AwaitingApproval);
3670 }
3671
3672 #[tokio::test]
3673 async fn approval_resume_completes_run() {
3674 let mut engine = create_test_engine();
3675 engine.register(SingleApprovalWorkflow).unwrap();
3676
3677 let run = engine
3679 .run_handler("single-approval", TriggerKind::Manual, json!({}))
3680 .await
3681 .unwrap()
3682 .run;
3683 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3684
3685 engine
3687 .store()
3688 .update_run_status(run.id, RunStatus::Running)
3689 .await
3690 .unwrap();
3691
3692 let resumed = engine.resume_run(run.id).await.unwrap().run;
3694 assert_eq!(resumed.status.state, RunStatus::Completed);
3695
3696 let steps = engine.store().list_steps(run.id).await.unwrap();
3697 assert_eq!(steps.len(), 3); assert_eq!(steps[0].name, "build");
3699 assert_eq!(steps[0].status.state, StepStatus::Completed);
3700 assert_eq!(steps[1].name, "gate");
3701 assert_eq!(steps[1].kind, StepKind::Approval);
3702 assert_eq!(steps[1].status.state, StepStatus::Completed);
3703 assert_eq!(steps[2].name, "deploy");
3704 assert_eq!(steps[2].status.state, StepStatus::Completed);
3705 }
3706
3707 #[tokio::test]
3708 async fn double_approval_two_resumes() {
3709 let mut engine = create_test_engine();
3710 engine.register(DoubleApprovalWorkflow).unwrap();
3711
3712 let run = engine
3714 .run_handler("double-approval", TriggerKind::Manual, json!({}))
3715 .await
3716 .unwrap()
3717 .run;
3718 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3719
3720 let steps = engine.store().list_steps(run.id).await.unwrap();
3721 assert_eq!(steps.len(), 2); engine
3725 .store()
3726 .update_run_status(run.id, RunStatus::Running)
3727 .await
3728 .unwrap();
3729
3730 let resumed = engine.resume_run(run.id).await.unwrap().run;
3731 assert_eq!(resumed.status.state, RunStatus::AwaitingApproval);
3732
3733 let steps = engine.store().list_steps(run.id).await.unwrap();
3734 assert_eq!(steps.len(), 4); engine
3738 .store()
3739 .update_run_status(run.id, RunStatus::Running)
3740 .await
3741 .unwrap();
3742
3743 let final_run = engine.resume_run(run.id).await.unwrap().run;
3744 assert_eq!(final_run.status.state, RunStatus::Completed);
3745
3746 let steps = engine.store().list_steps(run.id).await.unwrap();
3747 assert_eq!(steps.len(), 5);
3748 assert_eq!(steps[0].name, "build");
3749 assert_eq!(steps[1].name, "staging-gate");
3750 assert_eq!(steps[2].name, "deploy-staging");
3751 assert_eq!(steps[3].name, "prod-gate");
3752 assert_eq!(steps[4].name, "deploy-prod");
3753
3754 for step in &steps {
3755 assert_eq!(step.status.state, StepStatus::Completed);
3756 }
3757 }
3758
3759 use ironflow_store::models::{NewStep, StepUpdate, step_trace_id};
3764
3765 async fn create_step_with_status(
3766 store: &Arc<dyn Store>,
3767 run_id: Uuid,
3768 name: &str,
3769 position: u32,
3770 status: StepStatus,
3771 ) -> ironflow_store::models::Step {
3772 let step = store
3773 .create_step(NewStep {
3774 run_id,
3775 trace_id: step_trace_id(run_id, name, position),
3776 name: name.to_string(),
3777 kind: StepKind::Shell,
3778 position,
3779 input: None,
3780 is_error_handler: false,
3781 })
3782 .await
3783 .unwrap();
3784
3785 match status {
3786 StepStatus::Pending => {}
3787 StepStatus::Running => {
3788 store
3789 .update_step(
3790 step.id,
3791 StepUpdate {
3792 status: Some(StepStatus::Running),
3793 ..StepUpdate::default()
3794 },
3795 )
3796 .await
3797 .unwrap();
3798 }
3799 StepStatus::Completed => {
3800 store
3801 .update_step(
3802 step.id,
3803 StepUpdate {
3804 status: Some(StepStatus::Running),
3805 ..StepUpdate::default()
3806 },
3807 )
3808 .await
3809 .unwrap();
3810 store
3811 .update_step(
3812 step.id,
3813 StepUpdate {
3814 status: Some(StepStatus::Completed),
3815 ..StepUpdate::default()
3816 },
3817 )
3818 .await
3819 .unwrap();
3820 }
3821 StepStatus::AwaitingApproval => {
3822 store
3823 .update_step(
3824 step.id,
3825 StepUpdate {
3826 status: Some(StepStatus::Running),
3827 ..StepUpdate::default()
3828 },
3829 )
3830 .await
3831 .unwrap();
3832 store
3833 .update_step(
3834 step.id,
3835 StepUpdate {
3836 status: Some(StepStatus::AwaitingApproval),
3837 ..StepUpdate::default()
3838 },
3839 )
3840 .await
3841 .unwrap();
3842 }
3843 _ => panic!("unsupported status for test helper: {status}"),
3844 }
3845
3846 store.get_step(step.id).await.unwrap().unwrap()
3847 }
3848
3849 #[tokio::test]
3850 async fn fail_orphaned_steps_marks_running_as_failed() {
3851 let engine = create_test_engine();
3852 let run = engine
3853 .store()
3854 .create_run(NewRun {
3855 created_by: None,
3856 workflow_name: "test".to_string(),
3857 trigger: TriggerKind::Manual,
3858 payload: json!({}),
3859 max_retries: 0,
3860 handler_version: None,
3861 labels: HashMap::new(),
3862 scheduled_at: None,
3863 idempotency_key: None,
3864 concurrency_key: None,
3865 priority: 0,
3866 concurrency_limits: Vec::new(),
3867 max_cost_usd: None,
3868 worker_tags: Vec::new(),
3869 })
3870 .await
3871 .unwrap()
3872 .into_run();
3873
3874 let step = create_step_with_status(
3875 engine.store(),
3876 run.id,
3877 "running-step",
3878 0,
3879 StepStatus::Running,
3880 )
3881 .await;
3882
3883 engine
3884 .fail_orphaned_steps(run.id, "parent run timed out")
3885 .await
3886 .unwrap();
3887
3888 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
3889 assert_eq!(updated.status.state, StepStatus::Failed);
3890 assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
3891 assert!(updated.completed_at.is_some());
3892 }
3893
3894 #[tokio::test]
3895 async fn fail_orphaned_steps_marks_pending_as_skipped() {
3896 let engine = create_test_engine();
3897 let run = engine
3898 .store()
3899 .create_run(NewRun {
3900 created_by: None,
3901 workflow_name: "test".to_string(),
3902 trigger: TriggerKind::Manual,
3903 payload: json!({}),
3904 max_retries: 0,
3905 handler_version: None,
3906 labels: HashMap::new(),
3907 scheduled_at: None,
3908 idempotency_key: None,
3909 concurrency_key: None,
3910 priority: 0,
3911 concurrency_limits: Vec::new(),
3912 max_cost_usd: None,
3913 worker_tags: Vec::new(),
3914 })
3915 .await
3916 .unwrap()
3917 .into_run();
3918
3919 let step = create_step_with_status(
3920 engine.store(),
3921 run.id,
3922 "pending-step",
3923 0,
3924 StepStatus::Pending,
3925 )
3926 .await;
3927
3928 engine
3929 .fail_orphaned_steps(run.id, "parent run timed out")
3930 .await
3931 .unwrap();
3932
3933 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
3934 assert_eq!(updated.status.state, StepStatus::Skipped);
3935 assert!(updated.error.is_none());
3936 assert!(updated.completed_at.is_some());
3937 }
3938
3939 #[tokio::test]
3940 async fn fail_orphaned_steps_marks_awaiting_approval_as_failed() {
3941 let engine = create_test_engine();
3942 let run = engine
3943 .store()
3944 .create_run(NewRun {
3945 created_by: None,
3946 workflow_name: "test".to_string(),
3947 trigger: TriggerKind::Manual,
3948 payload: json!({}),
3949 max_retries: 0,
3950 handler_version: None,
3951 labels: HashMap::new(),
3952 scheduled_at: None,
3953 idempotency_key: None,
3954 concurrency_key: None,
3955 priority: 0,
3956 concurrency_limits: Vec::new(),
3957 max_cost_usd: None,
3958 worker_tags: Vec::new(),
3959 })
3960 .await
3961 .unwrap()
3962 .into_run();
3963
3964 let step = create_step_with_status(
3965 engine.store(),
3966 run.id,
3967 "approval-step",
3968 0,
3969 StepStatus::AwaitingApproval,
3970 )
3971 .await;
3972
3973 engine
3974 .fail_orphaned_steps(run.id, "parent run timed out")
3975 .await
3976 .unwrap();
3977
3978 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
3979 assert_eq!(updated.status.state, StepStatus::Failed);
3980 assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
3981 assert!(updated.completed_at.is_some());
3982 }
3983
3984 #[tokio::test]
3985 async fn fail_orphaned_steps_skips_terminal_steps() {
3986 let engine = create_test_engine();
3987 let run = engine
3988 .store()
3989 .create_run(NewRun {
3990 created_by: None,
3991 workflow_name: "test".to_string(),
3992 trigger: TriggerKind::Manual,
3993 payload: json!({}),
3994 max_retries: 0,
3995 handler_version: None,
3996 labels: HashMap::new(),
3997 scheduled_at: None,
3998 idempotency_key: None,
3999 concurrency_key: None,
4000 priority: 0,
4001 concurrency_limits: Vec::new(),
4002 max_cost_usd: None,
4003 worker_tags: Vec::new(),
4004 })
4005 .await
4006 .unwrap()
4007 .into_run();
4008
4009 let completed_step =
4010 create_step_with_status(engine.store(), run.id, "done", 0, StepStatus::Completed).await;
4011 let running_step =
4012 create_step_with_status(engine.store(), run.id, "in-flight", 1, StepStatus::Running)
4013 .await;
4014
4015 engine
4016 .fail_orphaned_steps(run.id, "parent run timed out")
4017 .await
4018 .unwrap();
4019
4020 let completed = engine
4021 .store()
4022 .get_step(completed_step.id)
4023 .await
4024 .unwrap()
4025 .unwrap();
4026 assert_eq!(completed.status.state, StepStatus::Completed);
4027
4028 let failed = engine
4029 .store()
4030 .get_step(running_step.id)
4031 .await
4032 .unwrap()
4033 .unwrap();
4034 assert_eq!(failed.status.state, StepStatus::Failed);
4035 }
4036
4037 #[tokio::test]
4038 async fn fail_orphaned_steps_mixed_states() {
4039 let engine = create_test_engine();
4040 let run = engine
4041 .store()
4042 .create_run(NewRun {
4043 created_by: None,
4044 workflow_name: "test".to_string(),
4045 trigger: TriggerKind::Manual,
4046 payload: json!({}),
4047 max_retries: 0,
4048 handler_version: None,
4049 labels: HashMap::new(),
4050 scheduled_at: None,
4051 idempotency_key: None,
4052 concurrency_key: None,
4053 priority: 0,
4054 concurrency_limits: Vec::new(),
4055 max_cost_usd: None,
4056 worker_tags: Vec::new(),
4057 })
4058 .await
4059 .unwrap()
4060 .into_run();
4061
4062 let s_completed =
4063 create_step_with_status(engine.store(), run.id, "step-1", 0, StepStatus::Completed)
4064 .await;
4065 let s_running =
4066 create_step_with_status(engine.store(), run.id, "step-2", 1, StepStatus::Running).await;
4067 let s_pending =
4068 create_step_with_status(engine.store(), run.id, "step-3", 2, StepStatus::Pending).await;
4069
4070 engine.fail_orphaned_steps(run.id, "timeout").await.unwrap();
4071
4072 let r_completed = engine
4073 .store()
4074 .get_step(s_completed.id)
4075 .await
4076 .unwrap()
4077 .unwrap();
4078 assert_eq!(r_completed.status.state, StepStatus::Completed);
4079
4080 let r_running = engine
4081 .store()
4082 .get_step(s_running.id)
4083 .await
4084 .unwrap()
4085 .unwrap();
4086 assert_eq!(r_running.status.state, StepStatus::Failed);
4087 assert_eq!(r_running.error.as_deref(), Some("timeout"));
4088
4089 let r_pending = engine
4090 .store()
4091 .get_step(s_pending.id)
4092 .await
4093 .unwrap()
4094 .unwrap();
4095 assert_eq!(r_pending.status.state, StepStatus::Skipped);
4096 assert!(r_pending.error.is_none());
4097 }
4098
4099 #[tokio::test]
4100 async fn fail_orphaned_steps_no_steps_is_noop() {
4101 let engine = create_test_engine();
4102 let run = engine
4103 .store()
4104 .create_run(NewRun {
4105 created_by: None,
4106 workflow_name: "test".to_string(),
4107 trigger: TriggerKind::Manual,
4108 payload: json!({}),
4109 max_retries: 0,
4110 handler_version: None,
4111 labels: HashMap::new(),
4112 scheduled_at: None,
4113 idempotency_key: None,
4114 concurrency_key: None,
4115 priority: 0,
4116 concurrency_limits: Vec::new(),
4117 max_cost_usd: None,
4118 worker_tags: Vec::new(),
4119 })
4120 .await
4121 .unwrap()
4122 .into_run();
4123
4124 let result = engine.fail_orphaned_steps(run.id, "timeout").await;
4125 assert!(result.is_ok());
4126 }
4127
4128 #[tokio::test]
4129 async fn fail_orphaned_steps_preserves_existing_error() {
4130 let engine = create_test_engine();
4131 let run = engine
4132 .store()
4133 .create_run(NewRun {
4134 created_by: None,
4135 workflow_name: "test".to_string(),
4136 trigger: TriggerKind::Manual,
4137 payload: json!({}),
4138 max_retries: 0,
4139 handler_version: None,
4140 labels: HashMap::new(),
4141 scheduled_at: None,
4142 idempotency_key: None,
4143 concurrency_key: None,
4144 priority: 0,
4145 concurrency_limits: Vec::new(),
4146 max_cost_usd: None,
4147 worker_tags: Vec::new(),
4148 })
4149 .await
4150 .unwrap()
4151 .into_run();
4152
4153 let step_with_error = create_step_with_status(
4154 engine.store(),
4155 run.id,
4156 "already-errored",
4157 0,
4158 StepStatus::Running,
4159 )
4160 .await;
4161
4162 engine
4163 .store()
4164 .update_step(
4165 step_with_error.id,
4166 StepUpdate {
4167 error: Some("real error from provider".to_string()),
4168 ..StepUpdate::default()
4169 },
4170 )
4171 .await
4172 .unwrap();
4173
4174 let step_no_error = create_step_with_status(
4175 engine.store(),
4176 run.id,
4177 "no-error-yet",
4178 1,
4179 StepStatus::Running,
4180 )
4181 .await;
4182
4183 engine
4184 .fail_orphaned_steps(run.id, "parent run failed")
4185 .await
4186 .unwrap();
4187
4188 let updated_with = engine
4189 .store()
4190 .get_step(step_with_error.id)
4191 .await
4192 .unwrap()
4193 .unwrap();
4194 assert_eq!(updated_with.status.state, StepStatus::Failed);
4195 assert_eq!(
4196 updated_with.error.as_deref(),
4197 Some("real error from provider"),
4198 );
4199
4200 let updated_without = engine
4201 .store()
4202 .get_step(step_no_error.id)
4203 .await
4204 .unwrap()
4205 .unwrap();
4206 assert_eq!(updated_without.status.state, StepStatus::Failed);
4207 assert_eq!(updated_without.error.as_deref(), Some("parent run failed"),);
4208 }
4209}