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 tokio::sync::{Mutex as AsyncMutex, OwnedMutexGuard};
19use tracing::{error, info, warn};
20use uuid::Uuid;
21
22use ironflow_core::error::OperationError;
23#[cfg(feature = "prometheus")]
24use ironflow_core::metric_names::{
25 RUN_BUDGET_EXCEEDED_TOTAL, RUN_COST_USD, RUN_DURATION_SECONDS, RUNS_ACTIVE, RUNS_TOTAL,
26};
27use ironflow_core::provider::{AgentProvider, LABEL_ROOT_RUN_ID};
28use ironflow_store::error::StoreError;
29use ironflow_store::models::{
30 ConcurrencyLimit, LeaseUpdate, NewRun, NewSignal, ProviderKind, Run, RunActor, RunCreation,
31 RunFilter, RunStatus, RunUpdate, SignalInsert, SignalStepResolution, StepStatus, StepUpdate,
32 TriggerKind, normalize_worker_tags, validate_concurrency_limits, validate_priority,
33 validate_worker_tags,
34};
35use ironflow_store::store::Store;
36#[cfg(feature = "prometheus")]
37use metrics::{counter, gauge, histogram};
38
39use crate::artifact::ArtifactSink;
40use crate::budget::{BudgetConfig, month_start};
41use crate::context::{PARENT_RUN_ID_LABEL, WorkflowContext, interrupt_running_steps};
42use crate::error::EngineError;
43use crate::executor::{StepInterceptor, StepResult};
44use crate::guard::{WorkflowGuardConfig, new_shared_guard_state};
45use crate::handler::{WorkflowHandler, WorkflowInfo, clamp_priority};
46use crate::log_sender::LogSender;
47use crate::notify::{
48 ApprovalRequestedEvent, Event, EventPublisher, EventSubscriber, RunBudgetExceededEvent,
49 RunFailedEvent, RunStatusChangedEvent, SignalAwaitedEvent, SignalReceivedEvent,
50 WorkflowEventBus,
51};
52use crate::plan::{
53 ExecutionPlan, PlanOptions, PlanRecorder, SharedPlanRecorder, estimate_durations, lock_plan,
54};
55use crate::retry_policy::{backoff_for_retry, is_run_retryable};
56use crate::schedule::CronSchedule;
57use crate::signal::{
58 Signal, SignalDelivery, SignalRejected, SignalResumed, received_output, validate_step_payload,
59};
60use ironflow_core::decision::DecisionProvider;
61
62#[derive(Debug, Clone)]
81pub struct WorkflowResult {
82 pub run: Run,
84 pub steps: Vec<StepResult>,
86}
87
88#[derive(Debug, Clone, Default)]
107pub struct EnqueueOptions {
108 pub max_retries: u32,
110 pub labels: HashMap<String, String>,
112 pub scheduled_at: Option<DateTime<Utc>>,
115 pub max_cost_usd: Option<Decimal>,
119 pub created_by: Option<RunActor>,
122 pub idempotency_key: Option<String>,
128 pub concurrency_key: Option<String>,
134 pub concurrency_limits: Vec<ConcurrencyLimit>,
142 pub priority: Option<i16>,
151 pub worker_tags: Vec<String>,
157}
158
159#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
170pub enum ExecutionMode {
171 #[default]
175 Local,
176 Workers,
179}
180
181pub struct Engine {
222 store: Arc<dyn Store>,
223 provider: Arc<dyn AgentProvider>,
224 handlers: HashMap<String, Arc<dyn WorkflowHandler>>,
225 event_publisher: EventPublisher,
226 log_sender: Option<LogSender>,
227 budget: BudgetConfig,
228 artifact_sink: Option<Arc<dyn ArtifactSink>>,
229 guard_config: Option<WorkflowGuardConfig>,
230 event_bus: Option<WorkflowEventBus>,
231 decision_provider: Option<Arc<dyn DecisionProvider>>,
232 step_interceptor: Option<Arc<dyn StepInterceptor>>,
233 execution_mode: ExecutionMode,
234 worker_tags: Option<Arc<Vec<String>>>,
235 active_runs: ActiveRuns,
236}
237
238type ActiveRuns = Arc<Mutex<HashMap<Uuid, Arc<AsyncMutex<()>>>>>;
240
241struct ActiveRunGuard {
246 runs: ActiveRuns,
247 run_id: Uuid,
248 entry: Arc<AsyncMutex<()>>,
249 lock: Option<OwnedMutexGuard<()>>,
250}
251
252impl Drop for ActiveRunGuard {
253 fn drop(&mut self) {
254 self.lock.take();
255 prune_active_run(&self.runs, self.run_id, &self.entry);
256 }
257}
258
259fn prune_active_run(runs: &ActiveRuns, run_id: Uuid, entry: &Arc<AsyncMutex<()>>) {
264 let mut runs = runs.lock().unwrap_or_else(|poisoned| poisoned.into_inner());
265 let unused = runs
266 .get(&run_id)
267 .is_some_and(|current| Arc::ptr_eq(current, entry))
268 && Arc::strong_count(entry) == 2;
269 if unused {
270 runs.remove(&run_id);
271 }
272}
273
274fn validate_category(handler_name: &str, category: &str) -> Result<(), EngineError> {
284 let reject = |reason: &str| {
285 Err(EngineError::InvalidWorkflow(format!(
286 "handler '{handler_name}' has invalid category '{category}': {reason}"
287 )))
288 };
289
290 if category.is_empty() {
291 return reject("empty category");
292 }
293 if category.starts_with('/') {
294 return reject("leading '/'");
295 }
296 if category.ends_with('/') {
297 return reject("trailing '/'");
298 }
299 for segment in category.split('/') {
300 if segment.is_empty() {
301 return reject("empty segment (double '/')");
302 }
303 if segment.trim().is_empty() {
304 return reject("whitespace-only segment");
305 }
306 }
307 Ok(())
308}
309
310fn chain_label(run: &Run, key: &str) -> Option<Uuid> {
315 if !matches!(run.trigger, TriggerKind::Workflow) {
316 return None;
317 }
318 let id = Uuid::parse_str(run.labels.get(key)?).ok()?;
319 (id != run.id).then_some(id)
320}
321
322pub fn chain_root(run: &Run) -> Option<Uuid> {
342 chain_label(run, LABEL_ROOT_RUN_ID)
343}
344
345fn chain_parent(run: &Run) -> Option<Uuid> {
347 chain_label(run, PARENT_RUN_ID_LABEL)
348}
349
350impl Engine {
351 pub fn new(store: Arc<dyn Store>, provider: Arc<dyn AgentProvider>) -> Self {
367 Self {
368 store,
369 provider,
370 handlers: HashMap::new(),
371 event_publisher: EventPublisher::new(),
372 log_sender: None,
373 budget: BudgetConfig::new(),
374 artifact_sink: None,
375 guard_config: None,
376 event_bus: None,
377 decision_provider: None,
378 step_interceptor: None,
379 execution_mode: ExecutionMode::default(),
380 worker_tags: None,
381 active_runs: Arc::new(Mutex::new(HashMap::new())),
382 }
383 }
384
385 async fn track_execution(&self, run_id: Uuid) -> ActiveRunGuard {
388 let entry = {
389 let mut runs = self
390 .active_runs
391 .lock()
392 .unwrap_or_else(|poisoned| poisoned.into_inner());
393 Arc::clone(runs.entry(run_id).or_default())
394 };
395 let lock = Arc::clone(&entry).lock_owned().await;
396 ActiveRunGuard {
397 runs: Arc::clone(&self.active_runs),
398 run_id,
399 entry,
400 lock: Some(lock),
401 }
402 }
403
404 async fn wait_until_idle(&self, run_id: Uuid) {
406 let entry = {
407 let runs = self
408 .active_runs
409 .lock()
410 .unwrap_or_else(|poisoned| poisoned.into_inner());
411 runs.get(&run_id).map(Arc::clone)
412 };
413 if let Some(entry) = entry {
414 drop(Arc::clone(&entry).lock_owned().await);
415 prune_active_run(&self.active_runs, run_id, &entry);
416 }
417 }
418
419 pub(crate) fn is_executing(&self, run_id: Uuid) -> bool {
421 let runs = self
422 .active_runs
423 .lock()
424 .unwrap_or_else(|poisoned| poisoned.into_inner());
425 runs.get(&run_id)
426 .is_some_and(|entry| entry.try_lock().is_err())
427 }
428
429 pub fn with_decision_provider(mut self, provider: Arc<dyn DecisionProvider>) -> Self {
450 self.decision_provider = Some(provider);
451 self
452 }
453
454 pub fn with_step_interceptor(mut self, interceptor: Arc<dyn StepInterceptor>) -> Self {
478 self.step_interceptor = Some(interceptor);
479 self
480 }
481
482 pub fn step_interceptor(&self) -> Option<&Arc<dyn StepInterceptor>> {
484 self.step_interceptor.as_ref()
485 }
486
487 pub fn with_budget_config(mut self, budget: BudgetConfig) -> Self {
508 self.budget = budget;
509 self
510 }
511
512 pub fn budget_config(&self) -> &BudgetConfig {
514 &self.budget
515 }
516
517 pub fn with_guard_config(mut self, config: WorkflowGuardConfig) -> Self {
539 self.guard_config = Some(config);
540 self
541 }
542
543 pub fn guard_config(&self) -> Option<&WorkflowGuardConfig> {
545 self.guard_config.as_ref()
546 }
547
548 pub fn with_execution_mode(mut self, mode: ExecutionMode) -> Self {
569 self.execution_mode = mode;
570 self
571 }
572
573 pub fn execution_mode(&self) -> ExecutionMode {
575 self.execution_mode
576 }
577
578 pub fn set_worker_tags(&mut self, tags: Vec<String>) {
601 self.worker_tags = Some(Arc::new(tags));
602 }
603
604 pub fn worker_tags(&self) -> Option<&[String]> {
622 self.worker_tags.as_deref().map(Vec::as_slice)
623 }
624
625 pub fn set_log_sender(&mut self, sender: LogSender) {
631 self.log_sender = Some(sender);
632 }
633
634 pub fn set_artifact_sink(&mut self, sink: Arc<dyn ArtifactSink>) {
652 self.artifact_sink = Some(sink);
653 }
654
655 pub fn artifact_sink(&self) -> Option<&Arc<dyn ArtifactSink>> {
657 self.artifact_sink.as_ref()
658 }
659
660 pub fn set_event_bus(&mut self, bus: WorkflowEventBus) {
677 self.event_bus = Some(bus);
678 }
679
680 pub fn event_bus(&self) -> Option<&WorkflowEventBus> {
682 self.event_bus.as_ref()
683 }
684
685 pub fn store(&self) -> &Arc<dyn Store> {
687 &self.store
688 }
689
690 pub fn provider(&self) -> &Arc<dyn AgentProvider> {
692 &self.provider
693 }
694
695 fn build_context(&self, run: &Run) -> WorkflowContext {
704 let handlers = self.handlers.clone();
705 let resolver: crate::context::HandlerResolver =
706 Arc::new(move |name: &str| handlers.get(name).cloned());
707 let mut ctx = WorkflowContext::with_handler_resolver(
708 run.id,
709 run.workflow_name.clone(),
710 self.store.clone(),
711 self.provider.clone(),
712 resolver,
713 );
714 ctx.carry_over_run_totals(run.retry_count + 1, run.cost_usd, run.duration_ms);
715 ctx.set_max_cost_usd(run.max_cost_usd);
716 ctx.set_run_created_at(run.created_at);
717 if let Some(ref sender) = self.log_sender {
718 ctx.set_log_sender(sender.clone());
719 }
720 if let Some(ref sink) = self.artifact_sink {
721 ctx.set_artifact_sink(sink.clone());
722 }
723 if let Some(ref bus) = self.event_bus {
724 ctx.set_event_bus(bus.clone());
725 }
726 if let Some(ref provider) = self.decision_provider {
727 ctx.set_decision_provider(provider.clone());
728 }
729 if let Some(ref interceptor) = self.step_interceptor {
730 ctx.set_step_interceptor(interceptor.clone());
731 }
732 if let Some(ref tags) = self.worker_tags {
733 ctx.set_worker_tags(tags.clone());
734 }
735 ctx
736 }
737
738 fn build_context_with_guard(
744 &self,
745 run: &Run,
746 handler: &dyn WorkflowHandler,
747 ) -> WorkflowContext {
748 let mut ctx = self.build_context(run);
749 let guard_config = handler.guard_config().or_else(|| self.guard_config.clone());
750 if let Some(config) = guard_config {
751 ctx.set_guard(config, new_shared_guard_state());
752 }
753 ctx
754 }
755
756 async fn check_monthly_quota(&self, workflow_name: &str) -> Result<(), EngineError> {
767 let Some(limit) = self.budget.monthly_cost_limit_usd else {
768 return Ok(());
769 };
770
771 let stats = self
772 .store
773 .get_stats(RunFilter {
774 created_after: Some(month_start(Utc::now())),
775 ..RunFilter::default()
776 })
777 .await?;
778
779 if stats.total_cost_usd < limit {
780 return Ok(());
781 }
782
783 warn!(
784 workflow = %workflow_name,
785 limit_usd = %limit,
786 spent_usd = %stats.total_cost_usd,
787 "monthly cost quota exhausted, refusing new run"
788 );
789
790 #[cfg(feature = "prometheus")]
791 counter!(
792 RUN_BUDGET_EXCEEDED_TOTAL,
793 "workflow" => workflow_name.to_string(),
794 "scope" => "monthly",
795 )
796 .increment(1);
797
798 Err(EngineError::MonthlyBudgetExceeded {
799 limit_usd: limit,
800 spent_usd: stats.total_cost_usd,
801 })
802 }
803
804 pub fn register(&mut self, handler: impl WorkflowHandler + 'static) -> Result<(), EngineError> {
848 let name = handler.name().to_string();
849 if self.handlers.contains_key(&name) {
850 return Err(EngineError::InvalidWorkflow(format!(
851 "handler '{}' already registered",
852 name
853 )));
854 }
855 if let Some(category) = handler.category() {
856 validate_category(&name, category)?;
857 }
858 self.handlers.insert(name, Arc::new(handler));
859 Ok(())
860 }
861
862 pub fn register_boxed(&mut self, handler: Box<dyn WorkflowHandler>) -> Result<(), EngineError> {
869 let name = handler.name().to_string();
870 if self.handlers.contains_key(&name) {
871 return Err(EngineError::InvalidWorkflow(format!(
872 "handler '{}' already registered",
873 name
874 )));
875 }
876 if let Some(category) = handler.category() {
877 validate_category(&name, category)?;
878 }
879 self.handlers.insert(name, Arc::from(handler));
880 Ok(())
881 }
882
883 pub fn get_handler(&self, name: &str) -> Option<&Arc<dyn WorkflowHandler>> {
885 self.handlers.get(name)
886 }
887
888 pub fn handler_names(&self) -> Vec<&str> {
890 self.handlers.keys().map(|s| s.as_str()).collect()
891 }
892
893 pub fn handler_info(&self, name: &str) -> Option<WorkflowInfo> {
895 self.handlers.get(name).map(|h| h.describe())
896 }
897
898 pub fn scheduled_handlers(&self) -> Vec<(&str, &CronSchedule)> {
922 self.handlers
923 .iter()
924 .filter_map(|(name, handler)| handler.schedule().map(|sched| (name.as_str(), sched)))
925 .collect()
926 }
927
928 pub fn subscribe(
953 &mut self,
954 subscriber: impl EventSubscriber + 'static,
955 event_types: &[&'static str],
956 ) {
957 self.event_publisher.subscribe(subscriber, event_types);
958 }
959
960 pub fn event_publisher(&self) -> &EventPublisher {
965 &self.event_publisher
966 }
967
968 #[tracing::instrument(name = "engine.run_handler", skip_all, fields(workflow = %handler_name))]
1004 pub async fn run_handler(
1005 &self,
1006 handler_name: &str,
1007 trigger: TriggerKind,
1008 payload: Value,
1009 ) -> Result<WorkflowResult, EngineError> {
1010 let handler = self
1011 .handlers
1012 .get(handler_name)
1013 .ok_or_else(|| {
1014 EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
1015 })?
1016 .clone();
1017
1018 self.check_monthly_quota(handler_name).await?;
1019
1020 let handler_version = handler.version().map(str::to_string);
1021 let max_cost_usd = self
1022 .budget
1023 .resolve_run_cap(None, handler.default_max_cost_usd());
1024 let run = self
1025 .store
1026 .create_run(NewRun {
1027 created_by: None,
1028 workflow_name: handler_name.to_string(),
1029 trigger,
1030 payload,
1031 max_retries: 0,
1032 handler_version,
1033 labels: handler.default_labels(),
1034 scheduled_at: None,
1035 idempotency_key: None,
1036 concurrency_key: None,
1037 priority: clamp_priority(handler.priority()),
1038 concurrency_limits: Vec::new(),
1039 max_cost_usd,
1040 worker_tags: normalize_worker_tags(handler.required_worker_tags()),
1041 })
1042 .await?
1043 .into_run();
1044
1045 let run_id = run.id;
1046 info!(run_id = %run_id, handler_version = run.handler_version.as_deref().unwrap_or(""), "run created");
1047
1048 if self.is_workflow_paused(handler_name).await? {
1049 info!(run_id = %run_id, workflow = %handler_name, "workflow paused, run left pending");
1050 return Ok(WorkflowResult {
1051 run,
1052 steps: Vec::new(),
1053 });
1054 }
1055
1056 self.store
1057 .update_run_status(run_id, RunStatus::Running)
1058 .await?;
1059
1060 #[cfg(feature = "prometheus")]
1061 gauge!(RUNS_ACTIVE, "workflow" => handler_name.to_string()).increment(1.0);
1062
1063 let run_start = Instant::now();
1064 let mut ctx = self.build_context_with_guard(&run, handler.as_ref());
1065
1066 let result = handler.execute(&mut ctx).await;
1067 self.finalize_run(run_id, handler_name, result, &ctx, run_start, run.labels)
1068 .await
1069 }
1070
1071 #[tracing::instrument(name = "engine.plan_handler", skip_all, fields(workflow = %handler_name))]
1113 pub async fn plan_handler(
1114 &self,
1115 handler_name: &str,
1116 payload: Value,
1117 options: PlanOptions,
1118 ) -> Result<ExecutionPlan, EngineError> {
1119 if options.max_depth == 0 {
1120 return Err(EngineError::InvalidWorkflow(
1121 "max_depth must be at least 1".to_string(),
1122 ));
1123 }
1124
1125 let handler = self
1126 .handlers
1127 .get(handler_name)
1128 .ok_or_else(|| {
1129 EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
1130 })?
1131 .clone();
1132
1133 let estimates = if options.estimate_durations {
1134 estimate_durations(&self.store, handler_name, options.sample_runs).await?
1135 } else {
1136 HashMap::new()
1137 };
1138
1139 let shared: SharedPlanRecorder = Arc::new(Mutex::new(PlanRecorder::new(
1140 handler_name.to_string(),
1141 payload,
1142 options.max_depth,
1143 estimates,
1144 )));
1145
1146 let handlers = self.handlers.clone();
1149 let resolver: crate::context::HandlerResolver =
1150 Arc::new(move |name: &str| handlers.get(name).cloned());
1151 let mut ctx = WorkflowContext::with_handler_resolver(
1152 Uuid::now_v7(),
1153 handler_name.to_string(),
1154 self.store.clone(),
1155 self.provider.clone(),
1156 resolver,
1157 );
1158 ctx.set_plan(shared.clone());
1159
1160 if let Err(err) = handler.execute(&mut ctx).await {
1161 lock_plan(&shared).fail(err.to_string());
1162 }
1163 drop(ctx);
1164
1165 let plan = match Arc::try_unwrap(shared) {
1166 Ok(mutex) => mutex
1167 .into_inner()
1168 .unwrap_or_else(|poisoned| poisoned.into_inner())
1169 .into_plan(),
1170 Err(shared) => lock_plan(&shared).snapshot(),
1171 };
1172
1173 info!(
1174 workflow = %handler_name,
1175 steps = plan.steps.len(),
1176 truncated = plan.truncated,
1177 "execution plan built"
1178 );
1179
1180 Ok(plan)
1181 }
1182
1183 #[tracing::instrument(name = "engine.enqueue_handler", skip_all, fields(workflow = %handler_name))]
1194 pub async fn enqueue_handler(
1195 &self,
1196 handler_name: &str,
1197 trigger: TriggerKind,
1198 payload: Value,
1199 max_retries: u32,
1200 ) -> Result<Run, EngineError> {
1201 self.enqueue_handler_with_options(
1202 handler_name,
1203 trigger,
1204 payload,
1205 EnqueueOptions {
1206 max_retries,
1207 ..Default::default()
1208 },
1209 )
1210 .await
1211 .map(RunCreation::into_run)
1212 }
1213
1214 #[tracing::instrument(name = "engine.enqueue_handler_with_options", skip_all, fields(workflow = %handler_name))]
1266 pub async fn enqueue_handler_with_options(
1267 &self,
1268 handler_name: &str,
1269 trigger: TriggerKind,
1270 payload: Value,
1271 options: EnqueueOptions,
1272 ) -> Result<RunCreation, EngineError> {
1273 let EnqueueOptions {
1274 max_retries,
1275 labels,
1276 scheduled_at,
1277 max_cost_usd,
1278 created_by,
1279 idempotency_key,
1280 concurrency_key,
1281 concurrency_limits,
1282 priority,
1283 worker_tags,
1284 } = options;
1285
1286 validate_concurrency_limits(&concurrency_limits)
1289 .map_err(EngineError::InvalidConcurrencyLimit)?;
1290 if let Some(priority) = priority {
1291 validate_priority(priority).map_err(EngineError::InvalidPriority)?;
1292 }
1293 validate_worker_tags(&worker_tags).map_err(EngineError::InvalidWorkerTag)?;
1294
1295 let handler = self.handlers.get(handler_name).ok_or_else(|| {
1296 EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
1297 })?;
1298
1299 self.check_monthly_quota(handler_name).await?;
1300
1301 let handler_version = handler.version().map(str::to_string);
1302 let mut merged_labels = handler.default_labels();
1303 merged_labels.extend(labels);
1304 let resolved_cap = self
1305 .budget
1306 .resolve_run_cap(max_cost_usd, handler.default_max_cost_usd());
1307 let priority = priority.unwrap_or_else(|| clamp_priority(handler.priority()));
1308 let required_tags = normalize_worker_tags(
1309 handler
1310 .required_worker_tags()
1311 .into_iter()
1312 .chain(worker_tags),
1313 );
1314
1315 let creation = self
1316 .store
1317 .create_run(NewRun {
1318 workflow_name: handler_name.to_string(),
1319 trigger,
1320 payload,
1321 max_retries,
1322 handler_version,
1323 labels: merged_labels,
1324 scheduled_at,
1325 created_by,
1326 idempotency_key,
1327 concurrency_key,
1328 priority,
1329 concurrency_limits,
1330 max_cost_usd: resolved_cap,
1331 worker_tags: required_tags,
1332 })
1333 .await?;
1334
1335 match &creation {
1336 RunCreation::Created(run) => info!(
1337 run_id = %run.id,
1338 workflow = %handler_name,
1339 max_cost_usd = ?resolved_cap,
1340 "handler run enqueued"
1341 ),
1342 RunCreation::Existing(run) => info!(
1343 run_id = %run.id,
1344 workflow = %handler_name,
1345 "idempotent replay, nothing enqueued"
1346 ),
1347 }
1348
1349 Ok(creation)
1350 }
1351
1352 #[tracing::instrument(name = "engine.execute_handler_run", skip_all, fields(run_id = %run_id))]
1372 pub async fn execute_handler_run(&self, run_id: Uuid) -> Result<WorkflowResult, EngineError> {
1373 let run = self
1374 .store
1375 .get_run(run_id)
1376 .await?
1377 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
1378
1379 if let Some(root_run_id) = chain_root(&run) {
1380 return self.resume_chain(run, root_run_id).await;
1381 }
1382
1383 let _active = self.track_execution(run_id).await;
1384
1385 let handler = self
1386 .handlers
1387 .get(&run.workflow_name)
1388 .ok_or_else(|| {
1389 EngineError::InvalidWorkflow(format!(
1390 "no handler registered: {}",
1391 run.workflow_name
1392 ))
1393 })?
1394 .clone();
1395
1396 #[cfg(feature = "prometheus")]
1397 gauge!(RUNS_ACTIVE, "workflow" => run.workflow_name.clone()).increment(1.0);
1398
1399 let run_start = Instant::now();
1400 let mut ctx = self.build_context_with_guard(&run, handler.as_ref());
1401
1402 let result = if handler.is_version_compatible(run.handler_version.as_deref()) {
1415 ctx.load_replay_steps().await?;
1416 self.release_then_execute(run_id, handler.as_ref(), &mut ctx)
1417 .await
1418 } else {
1419 Err(EngineError::HandlerVersionMismatch {
1420 run_id,
1421 workflow_name: run.workflow_name.clone(),
1422 run_version: run
1423 .handler_version
1424 .clone()
1425 .unwrap_or_else(|| "unknown".to_string()),
1426 current_version: handler
1427 .version()
1428 .map(str::to_string)
1429 .unwrap_or_else(|| "unknown".to_string()),
1430 })
1431 };
1432
1433 self.finalize_run(
1434 run_id,
1435 &run.workflow_name,
1436 result,
1437 &ctx,
1438 run_start,
1439 run.labels,
1440 )
1441 .await
1442 }
1443
1444 #[tracing::instrument(name = "engine.execute_run", skip_all, fields(run_id = %run_id))]
1452 pub async fn execute_run(&self, run_id: Uuid) -> Result<WorkflowResult, EngineError> {
1453 self.execute_handler_run(run_id).await
1454 }
1455
1456 #[tracing::instrument(name = "engine.resume_run", skip_all, fields(run_id = %run_id))]
1485 pub async fn resume_run(&self, run_id: Uuid) -> Result<WorkflowResult, EngineError> {
1486 let run = self
1487 .store
1488 .get_run(run_id)
1489 .await?
1490 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
1491
1492 if let Some(root_run_id) = chain_root(&run) {
1493 return self.resume_chain(run, root_run_id).await;
1494 }
1495
1496 self.resume_loaded_run(run).await
1497 }
1498
1499 async fn resume_chain(
1514 &self,
1515 child: Run,
1516 root_run_id: Uuid,
1517 ) -> Result<WorkflowResult, EngineError> {
1518 let child_run_id = child.id;
1519 let lease = child.worker_id.zip(child.lease_expires_at);
1520 let root = self
1521 .store
1522 .get_run(root_run_id)
1523 .await?
1524 .ok_or(EngineError::Store(StoreError::RunNotFound(root_run_id)))?;
1525
1526 match root.status.state {
1527 RunStatus::AwaitingApproval | RunStatus::Pending => {
1528 self.move_root_to_running(root_run_id, lease.as_ref())
1529 .await?;
1530 }
1531 RunStatus::Sleeping => {
1532 self.store
1533 .update_run_status(root_run_id, RunStatus::Pending)
1534 .await?;
1535 self.move_root_to_running(root_run_id, lease.as_ref())
1536 .await?;
1537 }
1538 other => {
1539 let reason = format!(
1540 "cannot resume child run {child_run_id}: root run {root_run_id} is {other}"
1541 );
1542 if let Err(err) = self
1543 .fail_or_schedule_retry(child_run_id, &reason, false, None, None)
1544 .await
1545 {
1546 error!(
1547 run_id = %child_run_id,
1548 error = %err,
1549 "failed to fail a child run whose root cannot resume"
1550 );
1551 }
1552 return Err(EngineError::InvalidWorkflow(reason));
1553 }
1554 }
1555
1556 if lease.is_some() {
1557 self.store
1558 .update_run(
1559 child_run_id,
1560 RunUpdate {
1561 lease: Some(LeaseUpdate::Release),
1562 ..RunUpdate::default()
1563 },
1564 )
1565 .await?;
1566 }
1567
1568 info!(
1569 run_id = %child_run_id,
1570 root_run_id = %root_run_id,
1571 lease_transferred = lease.is_some(),
1572 "child run resumed through its root run"
1573 );
1574
1575 let root = self
1576 .store
1577 .get_run(root_run_id)
1578 .await?
1579 .ok_or(EngineError::Store(StoreError::RunNotFound(root_run_id)))?;
1580 self.resume_loaded_run(root).await
1581 }
1582
1583 async fn move_root_to_running(
1589 &self,
1590 root_run_id: Uuid,
1591 lease: Option<&(String, DateTime<Utc>)>,
1592 ) -> Result<(), EngineError> {
1593 match lease {
1594 Some((worker_id, expires_at)) => {
1595 self.store
1596 .update_run(
1597 root_run_id,
1598 RunUpdate {
1599 status: Some(RunStatus::Running),
1600 lease: Some(LeaseUpdate::Set {
1601 worker_id: worker_id.clone(),
1602 expires_at: *expires_at,
1603 }),
1604 ..RunUpdate::default()
1605 },
1606 )
1607 .await?;
1608 }
1609 None => {
1610 self.store
1611 .update_run_status(root_run_id, RunStatus::Running)
1612 .await?;
1613 }
1614 }
1615 Ok(())
1616 }
1617
1618 async fn resume_loaded_run(&self, run: Run) -> Result<WorkflowResult, EngineError> {
1620 let run_id = run.id;
1621 let _active = self.track_execution(run_id).await;
1622 let handler = self
1623 .handlers
1624 .get(&run.workflow_name)
1625 .ok_or_else(|| {
1626 EngineError::InvalidWorkflow(format!(
1627 "no handler registered: {}",
1628 run.workflow_name
1629 ))
1630 })?
1631 .clone();
1632
1633 info!(run_id = %run_id, workflow = %run.workflow_name, "resuming run after approval");
1634
1635 let run_start = Instant::now();
1636 let mut ctx = self.build_context_with_guard(&run, handler.as_ref());
1637
1638 let result = if handler.is_version_compatible(run.handler_version.as_deref()) {
1639 ctx.load_replay_steps().await?;
1640 self.release_then_execute(run_id, handler.as_ref(), &mut ctx)
1641 .await
1642 } else {
1643 Err(EngineError::HandlerVersionMismatch {
1644 run_id,
1645 workflow_name: run.workflow_name.clone(),
1646 run_version: run
1647 .handler_version
1648 .clone()
1649 .unwrap_or_else(|| "unknown".to_string()),
1650 current_version: handler
1651 .version()
1652 .map(str::to_string)
1653 .unwrap_or_else(|| "unknown".to_string()),
1654 })
1655 };
1656
1657 self.finalize_run(
1658 run_id,
1659 &run.workflow_name,
1660 result,
1661 &ctx,
1662 run_start,
1663 run.labels,
1664 )
1665 .await
1666 }
1667
1668 pub async fn deliver_signal(
1714 self: &Arc<Self>,
1715 signal: NewSignal,
1716 ) -> Result<SignalDelivery, EngineError> {
1717 if signal.name.trim().is_empty() {
1718 return Err(EngineError::InvalidSignal(
1719 "signal name must not be empty".to_string(),
1720 ));
1721 }
1722 if signal.key.trim().is_empty() {
1723 return Err(EngineError::InvalidSignal(
1724 "signal key must not be empty".to_string(),
1725 ));
1726 }
1727
1728 let stored = match self.store.insert_signal(signal).await? {
1729 SignalInsert::Created(stored) => stored,
1730 SignalInsert::Duplicate(existing) => {
1731 info!(
1732 signal_id = %existing.id,
1733 signal = %existing.name,
1734 key = %existing.key,
1735 "duplicate signal ignored"
1736 );
1737 return Ok(SignalDelivery {
1738 signal_id: existing.id,
1739 duplicate: true,
1740 resumed: Vec::new(),
1741 rejected: Vec::new(),
1742 });
1743 }
1744 };
1745
1746 let waiters = self
1747 .store
1748 .list_signal_waiters(&stored.name, &stored.key)
1749 .await?;
1750 let mut resumed = Vec::new();
1751 let mut rejected = Vec::new();
1752
1753 for step in waiters {
1754 if let Err(error) = validate_step_payload(step.input.as_ref(), &stored.payload) {
1755 rejected.push(SignalRejected {
1756 run_id: step.run_id,
1757 step_id: step.id,
1758 error,
1759 });
1760 continue;
1761 }
1762
1763 match self
1764 .store
1765 .resolve_signal_step(step.id, received_output(&stored))
1766 .await
1767 {
1768 Ok(SignalStepResolution::Resolved {
1769 run_id,
1770 run_resumed,
1771 }) => {
1772 resumed.push(SignalResumed {
1773 run_id,
1774 step_id: step.id,
1775 });
1776 if run_resumed && self.execution_mode == ExecutionMode::Local {
1777 self.spawn_local_resume(run_id);
1778 }
1779 }
1780 Ok(SignalStepResolution::NotWaiting { .. }) => {}
1782 Err(err) => {
1783 error!(
1784 run_id = %step.run_id,
1785 step_id = %step.id,
1786 error = %err,
1787 "failed to resolve a waiting signal step"
1788 );
1789 rejected.push(SignalRejected {
1790 run_id: step.run_id,
1791 step_id: step.id,
1792 error: err.to_string(),
1793 });
1794 }
1795 }
1796 }
1797
1798 info!(
1799 signal_id = %stored.id,
1800 signal = %stored.name,
1801 key = %stored.key,
1802 resumed = resumed.len(),
1803 rejected = rejected.len(),
1804 "signal received"
1805 );
1806 self.event_publisher
1807 .publish(Event::SignalReceived(SignalReceivedEvent {
1808 signal_id: stored.id,
1809 name: stored.name.clone(),
1810 key: stored.key.clone(),
1811 resumed_runs: resumed.iter().map(|r| r.run_id).collect(),
1812 at: stored.received_at,
1813 }));
1814
1815 Ok(SignalDelivery {
1816 signal_id: stored.id,
1817 duplicate: false,
1818 resumed,
1819 rejected,
1820 })
1821 }
1822
1823 pub async fn send_signal<S: Signal>(
1861 self: &Arc<Self>,
1862 signal: &S,
1863 key: &str,
1864 idempotency_id: Option<&str>,
1865 ) -> Result<SignalDelivery, EngineError> {
1866 let payload = to_value(signal)?;
1867 self.deliver_signal(NewSignal {
1868 name: S::NAME.to_string(),
1869 key: key.to_string(),
1870 payload,
1871 idempotency_id: idempotency_id.map(str::to_string),
1872 })
1873 .await
1874 }
1875
1876 pub(crate) fn spawn_local_resume(self: &Arc<Self>, run_id: Uuid) {
1891 let engine = Arc::clone(self);
1892 spawn(async move {
1893 engine.wait_until_idle(run_id).await;
1894 let run = match engine.store.get_run(run_id).await {
1895 Ok(Some(run)) => run,
1896 Ok(None) => {
1897 error!(run_id = %run_id, "run to restart not found");
1898 return;
1899 }
1900 Err(err) => {
1901 error!(run_id = %run_id, error = %err, "failed to load a run to restart");
1902 return;
1903 }
1904 };
1905 match run.status.state {
1906 RunStatus::Pending => {
1907 if let Err(err) = engine
1908 .store
1909 .update_run_status(run_id, RunStatus::Running)
1910 .await
1911 {
1912 error!(run_id = %run_id, error = %err, "failed to restart a woken run");
1913 return;
1914 }
1915 }
1916 RunStatus::Running => {}
1917 status => {
1918 info!(
1919 run_id = %run_id,
1920 status = %status,
1921 "run no longer waiting to restart, resume skipped"
1922 );
1923 return;
1924 }
1925 }
1926 if let Err(err) = engine.resume_run(run_id).await {
1927 error!(run_id = %run_id, error = %err, "failed to resume a woken run");
1928 }
1929 });
1930 }
1931
1932 pub async fn fail_or_schedule_retry(
1983 &self,
1984 run_id: Uuid,
1985 error: &str,
1986 retryable: bool,
1987 cost_usd: Option<Decimal>,
1988 duration_ms: Option<u64>,
1989 ) -> Result<RunStatus, EngineError> {
1990 let run = self
1991 .store
1992 .get_run(run_id)
1993 .await?
1994 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
1995
1996 if run.status.state == RunStatus::Paused {
1999 info!(run_id = %run_id, error = %error, "run paused, failure not recorded");
2000 return Ok(RunStatus::Paused);
2001 }
2002
2003 let has_attempts_left = run.retry_count < run.max_retries;
2004 let update = if retryable && has_attempts_left {
2005 let backoff = backoff_for_retry(run.retry_count);
2006 let scheduled_at = Utc::now() + TimeDelta::milliseconds(backoff.as_millis() as i64);
2007
2008 info!(
2009 run_id = %run_id,
2010 workflow = %run.workflow_name,
2011 attempt = run.retry_count + 1,
2012 max_retries = run.max_retries,
2013 backoff_secs = backoff.as_secs(),
2014 scheduled_at = %scheduled_at,
2015 "run failed, scheduling retry"
2016 );
2017
2018 RunUpdate {
2019 status: Some(RunStatus::Retrying),
2020 error: Some(error.to_string()),
2021 increment_retry: true,
2022 cost_usd,
2023 duration_ms,
2024 scheduled_at: Some(scheduled_at),
2025 ..RunUpdate::default()
2026 }
2027 } else {
2028 RunUpdate {
2029 status: Some(RunStatus::Failed),
2030 error: Some(error.to_string()),
2031 cost_usd,
2032 duration_ms,
2033 completed_at: Some(Utc::now()),
2034 ..RunUpdate::default()
2035 }
2036 };
2037
2038 let status = update.status.unwrap_or(RunStatus::Failed);
2039 self.store.update_run(run_id, update).await?;
2040 self.fail_orphaned_steps(run_id, error).await?;
2041 self.cancel_descendants_of_stopped_run(run_id, error).await;
2044
2045 Ok(status)
2046 }
2047
2048 pub async fn interrupt_running_steps(&self, run_id: Uuid) -> Result<(), EngineError> {
2078 interrupt_running_steps(self.store.as_ref(), run_id).await
2079 }
2080
2081 pub async fn fail_orphaned_steps(
2095 &self,
2096 run_id: Uuid,
2097 error_message: &str,
2098 ) -> Result<(), EngineError> {
2099 let steps = self.store.list_steps(run_id).await?;
2100 let now = Utc::now();
2101
2102 for step in steps {
2103 if step.status.state.is_terminal() {
2104 continue;
2105 }
2106
2107 let (target_status, error) = match step.status.state {
2108 StepStatus::Running | StepStatus::AwaitingApproval => {
2109 let err = if step.error.is_some() {
2110 None
2111 } else {
2112 Some(error_message.to_string())
2113 };
2114 (StepStatus::Failed, err)
2115 }
2116 StepStatus::Pending => (StepStatus::Skipped, None),
2117 _ => continue,
2118 };
2119
2120 if let Err(e) = self
2121 .store
2122 .update_step(
2123 step.id,
2124 StepUpdate {
2125 status: Some(target_status),
2126 error,
2127 completed_at: Some(now),
2128 ..StepUpdate::default()
2129 },
2130 )
2131 .await
2132 {
2133 warn!(
2134 run_id = %run_id,
2135 step_id = %step.id,
2136 step_name = %step.name,
2137 error = %e,
2138 "failed to cleanup orphaned step"
2139 );
2140 } else {
2141 info!(
2142 run_id = %run_id,
2143 step_id = %step.id,
2144 step_name = %step.name,
2145 from = %step.status.state,
2146 to = %target_status,
2147 "cleaned up orphaned step"
2148 );
2149 }
2150 }
2151
2152 Ok(())
2153 }
2154
2155 async fn release_then_execute(
2160 &self,
2161 run_id: Uuid,
2162 handler: &dyn WorkflowHandler,
2163 ctx: &mut WorkflowContext,
2164 ) -> Result<(), EngineError> {
2165 match self.provider.release_run(&run_id.to_string()).await {
2166 Ok(()) => handler.execute(ctx).await,
2167 Err(e) => Err(EngineError::Operation(OperationError::Agent(e))),
2168 }
2169 }
2170
2171 async fn finalize_run(
2177 &self,
2178 run_id: Uuid,
2179 workflow_name: &str,
2180 result: Result<(), EngineError>,
2181 ctx: &WorkflowContext,
2182 run_start: Instant,
2183 run_labels: HashMap<String, String>,
2184 ) -> Result<WorkflowResult, EngineError> {
2185 let total_duration = ctx.carried_duration_ms() + run_start.elapsed().as_millis() as u64;
2188 let completed_at = Utc::now();
2189
2190 if let Some(run) = self.store.get_run(run_id).await?
2194 && run.status.state == RunStatus::Paused
2195 {
2196 let run = self
2197 .store
2198 .update_run_returning(
2199 run_id,
2200 RunUpdate {
2201 cost_usd: Some(ctx.total_cost_usd()),
2202 duration_ms: Some(total_duration),
2203 ..RunUpdate::default()
2204 },
2205 )
2206 .await?;
2207 info!(
2208 run_id = %run_id,
2209 outcome = ?result.err().map(|err| err.to_string()),
2210 "run paused, execution stopped"
2211 );
2212 return Ok(WorkflowResult {
2213 run,
2214 steps: ctx.step_results().to_vec(),
2215 });
2216 }
2217
2218 if matches!(result, Err(EngineError::RunPaused { .. })) {
2221 let run = self
2222 .store
2223 .get_run(run_id)
2224 .await?
2225 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
2226 info!(
2227 run_id = %run_id,
2228 status = %run.status.state,
2229 "run resumed after the pause, execution stopped"
2230 );
2231 return Ok(WorkflowResult {
2232 run,
2233 steps: ctx.step_results().to_vec(),
2234 });
2235 }
2236
2237 let final_status;
2238 let final_run;
2239
2240 match result {
2241 Ok(()) => {
2242 final_status = if ctx.has_allowed_failure() {
2243 RunStatus::Warning
2244 } else {
2245 RunStatus::Completed
2246 };
2247 final_run = self
2248 .store
2249 .update_run_returning(
2250 run_id,
2251 RunUpdate {
2252 status: Some(final_status),
2253 cost_usd: Some(ctx.total_cost_usd()),
2254 duration_ms: Some(total_duration),
2255 completed_at: Some(completed_at),
2256 output: ctx.output().cloned(),
2257 ..RunUpdate::default()
2258 },
2259 )
2260 .await?;
2261
2262 info!(
2263 run_id = %run_id,
2264 status = %final_status,
2265 cost_usd = %ctx.total_cost_usd(),
2266 duration_ms = total_duration,
2267 "run completed"
2268 );
2269 }
2270 Err(EngineError::ApprovalRequired {
2271 run_id: approval_run_id,
2272 step_id,
2273 ref message,
2274 }) => {
2275 final_status = RunStatus::AwaitingApproval;
2276 final_run = self
2277 .store
2278 .update_run_returning(
2279 run_id,
2280 RunUpdate {
2281 status: Some(RunStatus::AwaitingApproval),
2282 cost_usd: Some(ctx.total_cost_usd()),
2283 duration_ms: Some(total_duration),
2284 ..RunUpdate::default()
2285 },
2286 )
2287 .await?;
2288
2289 info!(
2290 run_id = %approval_run_id,
2291 step_id = %step_id,
2292 message = %message,
2293 "run awaiting approval"
2294 );
2295
2296 self.publish_approval_requested(approval_run_id, step_id, message)
2297 .await?;
2298 }
2299 Err(EngineError::ChildSuspended {
2300 run_id: child_run_id,
2301 ref cause,
2302 }) => {
2303 final_status = cause.suspension_status();
2304 final_run = self
2308 .store
2309 .update_run_returning(
2310 run_id,
2311 RunUpdate {
2312 status: Some(final_status),
2313 cost_usd: Some(ctx.total_cost_usd()),
2314 duration_ms: Some(total_duration),
2315 ..RunUpdate::default()
2316 },
2317 )
2318 .await?;
2319
2320 let leaf = cause.suspension_leaf();
2321 info!(
2322 run_id = %run_id,
2323 child_run_id = %child_run_id,
2324 status = %final_status,
2325 cause = %leaf,
2326 "run suspended with its child run"
2327 );
2328
2329 match leaf {
2330 EngineError::ApprovalRequired {
2331 run_id: approval_run_id,
2332 step_id,
2333 message,
2334 } => {
2335 self.publish_approval_requested(*approval_run_id, *step_id, message)
2336 .await?;
2337 }
2338 EngineError::SignalWaiting {
2339 run_id: wait_run_id,
2340 step_id,
2341 step_name,
2342 name,
2343 key,
2344 deadline_at,
2345 } => {
2346 self.event_publisher
2347 .publish(Event::SignalAwaited(SignalAwaitedEvent {
2348 run_id: *wait_run_id,
2349 step_id: *step_id,
2350 step_name: step_name.clone(),
2351 name: name.clone(),
2352 key: key.clone(),
2353 deadline_at: *deadline_at,
2354 at: Utc::now(),
2355 }));
2356 }
2357 _ => {}
2360 }
2361 }
2362 Err(EngineError::HumanInputRequired {
2363 run_id: input_run_id,
2364 step_id,
2365 ref message,
2366 }) => {
2367 final_status = RunStatus::AwaitingApproval;
2368 final_run = self
2369 .store
2370 .update_run_returning(
2371 run_id,
2372 RunUpdate {
2373 status: Some(RunStatus::AwaitingApproval),
2374 cost_usd: Some(ctx.total_cost_usd()),
2375 duration_ms: Some(total_duration),
2376 ..RunUpdate::default()
2377 },
2378 )
2379 .await?;
2380
2381 info!(
2383 run_id = %input_run_id,
2384 step_id = %step_id,
2385 message = %message,
2386 "run awaiting human input"
2387 );
2388 }
2389 Err(EngineError::DelaySleeping {
2390 run_id: delay_run_id,
2391 step_id,
2392 wake_at,
2393 }) => {
2394 final_status = RunStatus::Sleeping;
2395 final_run = self
2396 .store
2397 .update_run_returning(
2398 run_id,
2399 RunUpdate {
2400 status: Some(RunStatus::Sleeping),
2401 cost_usd: Some(ctx.total_cost_usd()),
2402 duration_ms: Some(total_duration),
2403 scheduled_at: Some(wake_at),
2404 ..RunUpdate::default()
2405 },
2406 )
2407 .await?;
2408
2409 info!(
2410 run_id = %delay_run_id,
2411 step_id = %step_id,
2412 wake_at = %wake_at,
2413 "run sleeping until delay elapses"
2414 );
2415 }
2416 Err(EngineError::CapacitySleeping {
2417 run_id: capacity_run_id,
2418 step_id,
2419 ref kind,
2420 wake_at,
2421 }) => {
2422 final_status = RunStatus::Sleeping;
2423 final_run = self
2424 .store
2425 .update_run_returning(
2426 run_id,
2427 RunUpdate {
2428 status: Some(RunStatus::Sleeping),
2429 cost_usd: Some(ctx.total_cost_usd()),
2430 duration_ms: Some(total_duration),
2431 scheduled_at: Some(wake_at),
2432 capacity_wait_kind: Some(ProviderKind::new(kind.as_str())),
2433 ..RunUpdate::default()
2434 },
2435 )
2436 .await?;
2437
2438 info!(
2439 run_id = %capacity_run_id,
2440 step_id = %step_id,
2441 kind = %kind,
2442 wake_at = %wake_at,
2443 "run sleeping until provider capacity returns"
2444 );
2445 }
2446 Err(EngineError::SignalWaiting {
2447 run_id: wait_run_id,
2448 step_id,
2449 ref step_name,
2450 ref name,
2451 ref key,
2452 deadline_at,
2453 }) => {
2454 final_status = RunStatus::Sleeping;
2455 let waiting = self
2459 .store
2460 .suspend_run_on_signal(run_id, step_id, deadline_at)
2461 .await?;
2462 final_run = self
2463 .store
2464 .update_run_returning(
2465 run_id,
2466 RunUpdate {
2467 cost_usd: Some(ctx.total_cost_usd()),
2468 duration_ms: Some(total_duration),
2469 ..RunUpdate::default()
2470 },
2471 )
2472 .await?;
2473
2474 if waiting {
2475 self.event_publisher
2476 .publish(Event::SignalAwaited(SignalAwaitedEvent {
2477 run_id: wait_run_id,
2478 step_id,
2479 step_name: step_name.clone(),
2480 name: name.clone(),
2481 key: key.clone(),
2482 deadline_at,
2483 at: Utc::now(),
2484 }));
2485 }
2486
2487 info!(
2488 run_id = %wait_run_id,
2489 step_id = %step_id,
2490 signal = %name,
2491 key = %key,
2492 deadline_at = %deadline_at,
2493 waiting,
2494 "run sleeping until a signal arrives"
2495 );
2496 }
2497 Err(err) => {
2498 let guardrail_stop = matches!(
2502 err,
2503 EngineError::RunBudgetExceeded { .. } | EngineError::WorkflowGuardRejected(_)
2504 );
2505
2506 final_status = if guardrail_stop {
2507 if let Err(store_err) = self
2508 .store
2509 .update_run(
2510 run_id,
2511 RunUpdate {
2512 status: Some(RunStatus::Cancelled),
2513 error: Some(err.to_string()),
2514 cost_usd: Some(ctx.total_cost_usd()),
2515 duration_ms: Some(total_duration),
2516 completed_at: Some(completed_at),
2517 output: ctx.output().cloned(),
2518 ..RunUpdate::default()
2519 },
2520 )
2521 .await
2522 {
2523 error!(run_id = %run_id, store_error = %store_err, "failed to persist run cancellation");
2524 }
2525 if let Err(cleanup_err) = self
2526 .fail_orphaned_steps(run_id, "run stopped: guardrail limit reached")
2527 .await
2528 {
2529 error!(run_id = %run_id, store_error = %cleanup_err, "failed to cleanup orphaned steps");
2530 }
2531 RunStatus::Cancelled
2532 } else {
2533 if let Some(output) = ctx.output()
2536 && let Err(store_err) = self
2537 .store
2538 .update_run(
2539 run_id,
2540 RunUpdate {
2541 output: Some(output.clone()),
2542 ..RunUpdate::default()
2543 },
2544 )
2545 .await
2546 {
2547 error!(run_id = %run_id, store_error = %store_err, "failed to persist run output");
2548 }
2549 self.fail_or_schedule_retry(
2550 run_id,
2551 &err.to_string(),
2552 is_run_retryable(&err),
2553 Some(ctx.total_cost_usd()),
2554 Some(total_duration),
2555 )
2556 .await
2557 .unwrap_or_else(|store_err| {
2558 error!(run_id = %run_id, store_error = %store_err, "failed to persist run failure");
2559 RunStatus::Failed
2560 })
2561 };
2562
2563 if matches!(err, EngineError::RunBudgetExceeded { .. }) {
2564 self.on_run_budget_exceeded(workflow_name, run_id, &err);
2565 }
2566
2567 error!(run_id = %run_id, status = %final_status, error = %err, "run stopped");
2568
2569 self.publish_run_status_changed(
2570 workflow_name,
2571 run_id,
2572 final_status,
2573 Some(err.to_string()),
2574 ctx,
2575 total_duration,
2576 run_labels,
2577 );
2578
2579 #[cfg(feature = "prometheus")]
2580 self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
2581
2582 return Err(err);
2583 }
2584 }
2585
2586 self.publish_run_status_changed(
2587 workflow_name,
2588 run_id,
2589 final_status,
2590 None,
2591 ctx,
2592 total_duration,
2593 run_labels,
2594 );
2595
2596 #[cfg(feature = "prometheus")]
2597 self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
2598
2599 Ok(WorkflowResult {
2600 run: final_run,
2601 steps: ctx.step_results().to_vec(),
2602 })
2603 }
2604
2605 async fn publish_approval_requested(
2608 &self,
2609 run_id: Uuid,
2610 step_id: Uuid,
2611 message: &str,
2612 ) -> Result<(), EngineError> {
2613 let requirement = self
2614 .store
2615 .get_step(step_id)
2616 .await?
2617 .and_then(|s| s.approval_requirement);
2618 self.event_publisher
2619 .publish(Event::ApprovalRequested(ApprovalRequestedEvent {
2620 run_id,
2621 step_id,
2622 message: message.to_string(),
2623 requirement,
2624 at: Utc::now(),
2625 }));
2626 Ok(())
2627 }
2628
2629 pub async fn fail_ancestors(&self, run_id: Uuid, reason: &str) -> Result<(), EngineError> {
2662 let mut current = self
2663 .store
2664 .get_run(run_id)
2665 .await?
2666 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
2667 let mut visited = HashSet::from([run_id]);
2669
2670 while let Some(parent_id) = chain_parent(¤t) {
2671 if !visited.insert(parent_id) {
2672 break;
2673 }
2674 let status = self
2675 .fail_or_schedule_retry(parent_id, reason, false, None, None)
2676 .await?;
2677 if status == RunStatus::Paused {
2680 self.requeue_paused_root(¤t).await?;
2681 break;
2682 }
2683 info!(
2684 run_id = %run_id,
2685 ancestor_run_id = %parent_id,
2686 status = %status,
2687 "ancestor run failed with its child"
2688 );
2689 current = self
2690 .store
2691 .get_run(parent_id)
2692 .await?
2693 .ok_or(EngineError::Store(StoreError::RunNotFound(parent_id)))?;
2694 }
2695
2696 Ok(())
2697 }
2698
2699 #[cfg(feature = "prometheus")]
2701 fn emit_run_metrics(
2702 &self,
2703 workflow_name: &str,
2704 status: RunStatus,
2705 duration_ms: u64,
2706 ctx: &WorkflowContext,
2707 ) {
2708 let status_str = status.to_string();
2709 let wf = workflow_name.to_string();
2710
2711 counter!(RUNS_TOTAL, "workflow" => wf.clone(), "status" => status_str.clone()).increment(1);
2712 histogram!(RUN_DURATION_SECONDS, "workflow" => wf.clone(), "status" => status_str)
2713 .record(duration_ms as f64 / 1000.0);
2714 histogram!(RUN_COST_USD, "workflow" => wf.clone()).record(
2715 ctx.total_cost_usd()
2716 .to_string()
2717 .parse::<f64>()
2718 .unwrap_or(0.0),
2719 );
2720 gauge!(RUNS_ACTIVE, "workflow" => wf).decrement(1.0);
2721 }
2722
2723 fn on_run_budget_exceeded(&self, workflow_name: &str, run_id: Uuid, err: &EngineError) {
2729 let EngineError::RunBudgetExceeded {
2730 limit_usd,
2731 spent_usd,
2732 step_budget_usd,
2733 ..
2734 } = err
2735 else {
2736 return;
2737 };
2738
2739 #[cfg(feature = "prometheus")]
2740 counter!(
2741 RUN_BUDGET_EXCEEDED_TOTAL,
2742 "workflow" => workflow_name.to_string(),
2743 "scope" => "run",
2744 )
2745 .increment(1);
2746
2747 self.event_publisher
2748 .publish(Event::RunBudgetExceeded(RunBudgetExceededEvent {
2749 run_id,
2750 workflow_name: workflow_name.to_string(),
2751 limit_usd: *limit_usd,
2752 spent_usd: *spent_usd,
2753 step_budget_usd: *step_budget_usd,
2754 at: Utc::now(),
2755 }));
2756 }
2757
2758 #[allow(clippy::too_many_arguments)]
2763 fn publish_run_status_changed(
2764 &self,
2765 workflow_name: &str,
2766 run_id: Uuid,
2767 to: RunStatus,
2768 error: Option<String>,
2769 ctx: &WorkflowContext,
2770 duration_ms: u64,
2771 labels: HashMap<String, String>,
2772 ) {
2773 let now = Utc::now();
2774 let cost_usd = ctx.total_cost_usd();
2775 let wf = workflow_name.to_string();
2776
2777 self.event_publisher
2778 .publish(Event::RunStatusChanged(RunStatusChangedEvent {
2779 run_id,
2780 workflow_name: wf.clone(),
2781 from: RunStatus::Running,
2782 to,
2783 error: error.clone(),
2784 cost_usd,
2785 duration_ms,
2786 labels: labels.clone(),
2787 at: now,
2788 }));
2789
2790 if to == RunStatus::Failed {
2791 self.event_publisher
2792 .publish(Event::RunFailed(RunFailedEvent {
2793 run_id,
2794 workflow_name: wf,
2795 error,
2796 cost_usd,
2797 duration_ms,
2798 labels,
2799 at: now,
2800 }));
2801 }
2802 }
2803}
2804
2805impl fmt::Debug for Engine {
2806 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
2807 f.debug_struct("Engine")
2808 .field("handlers", &self.handlers.keys().collect::<Vec<_>>())
2809 .finish_non_exhaustive()
2810 }
2811}
2812
2813#[cfg(test)]
2814mod tests {
2815 use super::*;
2816 use crate::config::ShellConfig;
2817 use crate::handler::{HandlerFuture, WorkflowHandler};
2818 use ironflow_core::providers::claude::ClaudeCodeProvider;
2819 use ironflow_core::providers::record_replay::RecordReplayProvider;
2820 use ironflow_store::memory::InMemoryStore;
2821 use ironflow_store::models::{MAX_PRIORITY, MIN_PRIORITY, StepStatus};
2822 use serde_json::json;
2823
2824 struct EchoWorkflow;
2826
2827 impl WorkflowHandler for EchoWorkflow {
2828 fn name(&self) -> &str {
2829 "echo-workflow"
2830 }
2831
2832 fn describe(&self) -> WorkflowInfo {
2833 WorkflowInfo {
2834 description: "A simple workflow that echoes hello".to_string(),
2835 source_code: None,
2836 sub_workflows: Vec::new(),
2837 category: None,
2838 version: self.version().map(str::to_string),
2839 compatible_versions: Vec::new(),
2840 input_schema: None,
2841 default_labels: HashMap::new(),
2842 schedule: self.schedule().cloned(),
2843 default_max_cost_usd: self.default_max_cost_usd(),
2844 priority: self.priority(),
2845 }
2846 }
2847
2848 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2849 Box::pin(async move {
2850 ctx.shell("greet", ShellConfig::new("echo hello")).await?;
2851 Ok(())
2852 })
2853 }
2854 }
2855
2856 struct FailingWorkflow;
2858
2859 impl WorkflowHandler for FailingWorkflow {
2860 fn name(&self) -> &str {
2861 "failing-workflow"
2862 }
2863
2864 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2865 Box::pin(async move {
2866 ctx.shell("fail", ShellConfig::new("exit 1")).await?;
2867 Ok(())
2868 })
2869 }
2870 }
2871
2872 fn create_test_engine() -> Engine {
2873 let store = Arc::new(InMemoryStore::new());
2874 let inner = ClaudeCodeProvider::new();
2875 let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
2876 inner,
2877 "/tmp/ironflow-fixtures",
2878 ));
2879 Engine::new(store, provider)
2880 }
2881
2882 #[test]
2883 fn engine_new_creates_instance() {
2884 let engine = create_test_engine();
2885 assert_eq!(engine.handler_names().len(), 0);
2886 }
2887
2888 #[test]
2889 fn execution_mode_defaults_to_local() {
2890 let engine = create_test_engine();
2891 assert_eq!(engine.execution_mode(), ExecutionMode::Local);
2892 }
2893
2894 #[test]
2895 fn with_execution_mode_overrides_the_default() {
2896 let engine = create_test_engine().with_execution_mode(ExecutionMode::Workers);
2897 assert_eq!(engine.execution_mode(), ExecutionMode::Workers);
2898 }
2899
2900 #[test]
2901 fn engine_register_handler() {
2902 let mut engine = create_test_engine();
2903 let result = engine.register(EchoWorkflow);
2904 assert!(result.is_ok());
2905 assert_eq!(engine.handler_names().len(), 1);
2906 assert!(engine.handler_names().contains(&"echo-workflow"));
2907 }
2908
2909 #[test]
2910 fn engine_register_duplicate_returns_error() {
2911 let mut engine = create_test_engine();
2912 engine.register(EchoWorkflow).unwrap();
2913 let result = engine.register(EchoWorkflow);
2914 assert!(result.is_err());
2915 }
2916
2917 #[test]
2918 fn engine_get_handler_found() {
2919 let mut engine = create_test_engine();
2920 engine.register(EchoWorkflow).unwrap();
2921 let handler = engine.get_handler("echo-workflow");
2922 assert!(handler.is_some());
2923 }
2924
2925 #[test]
2926 fn engine_get_handler_not_found() {
2927 let engine = create_test_engine();
2928 let handler = engine.get_handler("nonexistent");
2929 assert!(handler.is_none());
2930 }
2931
2932 #[test]
2933 fn engine_handler_names_lists_all() {
2934 let mut engine = create_test_engine();
2935 engine.register(EchoWorkflow).unwrap();
2936 engine.register(FailingWorkflow).unwrap();
2937 let names = engine.handler_names();
2938 assert_eq!(names.len(), 2);
2939 assert!(names.contains(&"echo-workflow"));
2940 assert!(names.contains(&"failing-workflow"));
2941 }
2942
2943 #[test]
2944 fn engine_handler_info_returns_description() {
2945 let mut engine = create_test_engine();
2946 engine.register(EchoWorkflow).unwrap();
2947 let info = engine.handler_info("echo-workflow");
2948 assert!(info.is_some());
2949 let info = info.unwrap();
2950 assert_eq!(info.description, "A simple workflow that echoes hello");
2951 }
2952
2953 struct CategorizedWorkflow;
2954
2955 impl WorkflowHandler for CategorizedWorkflow {
2956 fn name(&self) -> &str {
2957 "categorized"
2958 }
2959 fn category(&self) -> Option<&str> {
2960 Some("data/etl")
2961 }
2962 fn execute<'a>(
2963 &'a self,
2964 _ctx: &'a mut WorkflowContext,
2965 ) -> crate::handler::HandlerFuture<'a> {
2966 Box::pin(async move { Ok(()) })
2967 }
2968 }
2969
2970 #[test]
2971 fn engine_default_describe_propagates_category() {
2972 let mut engine = create_test_engine();
2973 engine.register(CategorizedWorkflow).unwrap();
2974 let info = engine.handler_info("categorized").unwrap();
2975 assert_eq!(info.category.as_deref(), Some("data/etl"));
2976 }
2977
2978 #[test]
2979 fn engine_default_describe_without_category() {
2980 let mut engine = create_test_engine();
2981 engine.register(EchoWorkflow).unwrap();
2982 let info = engine.handler_info("echo-workflow").unwrap();
2983 assert!(info.category.is_none());
2984 }
2985
2986 struct ScheduledWorkflow {
2991 schedule: CronSchedule,
2992 }
2993
2994 impl ScheduledWorkflow {
2995 fn new() -> Self {
2996 Self {
2997 schedule: CronSchedule::new("0 0 * * * *").unwrap(),
2998 }
2999 }
3000 }
3001
3002 impl WorkflowHandler for ScheduledWorkflow {
3003 fn name(&self) -> &str {
3004 "scheduled"
3005 }
3006 fn schedule(&self) -> Option<&CronSchedule> {
3007 Some(&self.schedule)
3008 }
3009 fn execute<'a>(
3010 &'a self,
3011 _ctx: &'a mut WorkflowContext,
3012 ) -> crate::handler::HandlerFuture<'a> {
3013 Box::pin(async move { Ok(()) })
3014 }
3015 }
3016
3017 #[test]
3018 fn engine_default_describe_propagates_schedule() {
3019 let mut engine = create_test_engine();
3020 engine.register(ScheduledWorkflow::new()).unwrap();
3021 let info = engine.handler_info("scheduled").unwrap();
3022 assert_eq!(
3023 info.schedule.as_ref().map(|s| s.as_str()),
3024 Some("0 0 * * * *")
3025 );
3026 }
3027
3028 #[test]
3029 fn engine_default_describe_without_schedule() {
3030 let mut engine = create_test_engine();
3031 engine.register(EchoWorkflow).unwrap();
3032 let info = engine.handler_info("echo-workflow").unwrap();
3033 assert!(info.schedule.is_none());
3034 }
3035
3036 #[test]
3037 fn scheduled_handlers_returns_only_scheduled() {
3038 let mut engine = create_test_engine();
3039 engine.register(EchoWorkflow).unwrap();
3040 engine.register(ScheduledWorkflow::new()).unwrap();
3041 engine.register(FailingWorkflow).unwrap();
3042
3043 let scheduled = engine.scheduled_handlers();
3044 assert_eq!(scheduled.len(), 1);
3045 assert_eq!(scheduled[0].0, "scheduled");
3046 assert_eq!(scheduled[0].1.as_str(), "0 0 * * * *");
3047 }
3048
3049 #[test]
3050 fn scheduled_handlers_empty_when_none_scheduled() {
3051 let mut engine = create_test_engine();
3052 engine.register(EchoWorkflow).unwrap();
3053 engine.register(FailingWorkflow).unwrap();
3054
3055 let scheduled = engine.scheduled_handlers();
3056 assert!(scheduled.is_empty());
3057 }
3058
3059 struct BadCategoryWorkflow(&'static str);
3060
3061 impl WorkflowHandler for BadCategoryWorkflow {
3062 fn name(&self) -> &str {
3063 "bad-category"
3064 }
3065 fn category(&self) -> Option<&str> {
3066 Some(self.0)
3067 }
3068 fn execute<'a>(
3069 &'a self,
3070 _ctx: &'a mut WorkflowContext,
3071 ) -> crate::handler::HandlerFuture<'a> {
3072 Box::pin(async move { Ok(()) })
3073 }
3074 }
3075
3076 #[test]
3077 fn engine_register_rejects_empty_category() {
3078 let mut engine = create_test_engine();
3079 let err = engine.register(BadCategoryWorkflow("")).unwrap_err();
3080 match err {
3081 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty category")),
3082 other => panic!("expected InvalidWorkflow, got {other:?}"),
3083 }
3084 }
3085
3086 #[test]
3087 fn engine_register_rejects_leading_slash_category() {
3088 let mut engine = create_test_engine();
3089 let err = engine
3090 .register(BadCategoryWorkflow("/data/etl"))
3091 .unwrap_err();
3092 match err {
3093 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("leading '/'")),
3094 other => panic!("expected InvalidWorkflow, got {other:?}"),
3095 }
3096 }
3097
3098 #[test]
3099 fn engine_register_rejects_trailing_slash_category() {
3100 let mut engine = create_test_engine();
3101 let err = engine
3102 .register(BadCategoryWorkflow("data/etl/"))
3103 .unwrap_err();
3104 match err {
3105 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("trailing '/'")),
3106 other => panic!("expected InvalidWorkflow, got {other:?}"),
3107 }
3108 }
3109
3110 #[test]
3111 fn engine_register_rejects_double_slash_category() {
3112 let mut engine = create_test_engine();
3113 let err = engine
3114 .register(BadCategoryWorkflow("data//etl"))
3115 .unwrap_err();
3116 match err {
3117 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty segment")),
3118 other => panic!("expected InvalidWorkflow, got {other:?}"),
3119 }
3120 }
3121
3122 #[test]
3123 fn engine_register_rejects_whitespace_only_segment_category() {
3124 let mut engine = create_test_engine();
3125 let err = engine
3126 .register(BadCategoryWorkflow("data/ /etl"))
3127 .unwrap_err();
3128 match err {
3129 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("whitespace-only segment")),
3130 other => panic!("expected InvalidWorkflow, got {other:?}"),
3131 }
3132 }
3133
3134 #[test]
3135 fn engine_register_accepts_valid_nested_category() {
3136 let mut engine = create_test_engine();
3137 assert!(engine.register(CategorizedWorkflow).is_ok());
3138 }
3139
3140 #[tokio::test]
3141 async fn engine_unknown_workflow_returns_error() {
3142 let engine = create_test_engine();
3143 let result = engine
3144 .run_handler("unknown", TriggerKind::Manual, json!({}))
3145 .await;
3146 assert!(result.is_err());
3147 match result {
3148 Err(EngineError::InvalidWorkflow(msg)) => {
3149 assert!(msg.contains("no handler registered"));
3150 }
3151 _ => panic!("expected InvalidWorkflow error"),
3152 }
3153 }
3154
3155 #[tokio::test]
3156 async fn engine_enqueue_handler_creates_pending_run() {
3157 let mut engine = create_test_engine();
3158 engine.register(EchoWorkflow).unwrap();
3159
3160 let run = engine
3161 .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
3162 .await
3163 .unwrap();
3164 assert_eq!(run.status.state, RunStatus::Pending);
3165 assert_eq!(run.workflow_name, "echo-workflow");
3166 }
3167
3168 #[tokio::test]
3169 async fn enqueue_handler_leaves_the_run_unattributed() {
3170 let mut engine = create_test_engine();
3171 engine.register(EchoWorkflow).unwrap();
3172
3173 let run = engine
3174 .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
3175 .await
3176 .unwrap();
3177
3178 assert!(run.created_by.is_none());
3179 }
3180
3181 #[tokio::test]
3182 async fn enqueue_handler_with_options_records_the_author() {
3183 let mut engine = create_test_engine();
3184 engine.register(EchoWorkflow).unwrap();
3185 let actor = RunActor::User {
3186 user_id: Uuid::now_v7(),
3187 };
3188
3189 let run = engine
3190 .enqueue_handler_with_options(
3191 "echo-workflow",
3192 TriggerKind::Api,
3193 json!({}),
3194 EnqueueOptions {
3195 created_by: Some(actor.clone()),
3196 ..Default::default()
3197 },
3198 )
3199 .await
3200 .unwrap()
3201 .into_run();
3202
3203 assert_eq!(run.created_by, Some(actor));
3204 }
3205
3206 #[tokio::test]
3207 async fn enqueue_handler_with_options_accepts_no_author() {
3208 let mut engine = create_test_engine();
3209 engine.register(EchoWorkflow).unwrap();
3210
3211 let run = engine
3212 .enqueue_handler_with_options(
3213 "echo-workflow",
3214 TriggerKind::Cron {
3215 schedule: "0 * * * * *".to_string(),
3216 schedule_id: None,
3217 scheduled_for: None,
3218 },
3219 json!({}),
3220 EnqueueOptions::default(),
3221 )
3222 .await
3223 .unwrap()
3224 .into_run();
3225
3226 assert!(run.created_by.is_none());
3227 }
3228
3229 #[tokio::test]
3230 async fn enqueue_handler_with_options_stores_concurrency_limits() {
3231 let mut engine = create_test_engine();
3232 engine.register(EchoWorkflow).unwrap();
3233 let limits = vec![
3234 ConcurrencyLimit::new("repo:acme", 2),
3235 ConcurrencyLimit::new("tenant:42", 5),
3236 ];
3237
3238 let run = engine
3239 .enqueue_handler_with_options(
3240 "echo-workflow",
3241 TriggerKind::Api,
3242 json!({}),
3243 EnqueueOptions {
3244 concurrency_limits: limits.clone(),
3245 ..Default::default()
3246 },
3247 )
3248 .await
3249 .unwrap()
3250 .into_run();
3251
3252 assert_eq!(run.concurrency_limits, limits);
3253 }
3254
3255 #[tokio::test]
3256 async fn enqueue_rejects_invalid_concurrency_limits() {
3257 let mut engine = create_test_engine();
3258 engine.register(EchoWorkflow).unwrap();
3259
3260 let invalid = [
3261 vec![ConcurrencyLimit::new("repo:acme", 0)],
3262 vec![ConcurrencyLimit::new("", 1)],
3263 vec![
3264 ConcurrencyLimit::new("repo:acme", 1),
3265 ConcurrencyLimit::new("repo:acme", 2),
3266 ],
3267 ];
3268 for concurrency_limits in invalid {
3269 let err = engine
3270 .enqueue_handler_with_options(
3271 "echo-workflow",
3272 TriggerKind::Api,
3273 json!({}),
3274 EnqueueOptions {
3275 concurrency_limits,
3276 ..Default::default()
3277 },
3278 )
3279 .await
3280 .unwrap_err();
3281 assert!(
3282 matches!(err, EngineError::InvalidConcurrencyLimit(_)),
3283 "{err:?}"
3284 );
3285 }
3286
3287 let err = engine
3289 .enqueue_handler_with_options(
3290 "not-registered",
3291 TriggerKind::Api,
3292 json!({}),
3293 EnqueueOptions {
3294 concurrency_limits: vec![ConcurrencyLimit::new("repo:acme", 0)],
3295 ..Default::default()
3296 },
3297 )
3298 .await
3299 .unwrap_err();
3300 assert!(
3301 matches!(err, EngineError::InvalidConcurrencyLimit(_)),
3302 "{err:?}"
3303 );
3304
3305 let page = engine
3306 .store()
3307 .list_runs(RunFilter::default(), 1, 10)
3308 .await
3309 .unwrap();
3310 assert_eq!(page.total, 0, "no run may be created");
3311 }
3312
3313 struct UrgentWorkflow;
3314
3315 impl WorkflowHandler for UrgentWorkflow {
3316 fn name(&self) -> &str {
3317 "urgent-workflow"
3318 }
3319
3320 fn priority(&self) -> i16 {
3321 60
3322 }
3323
3324 fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3325 Box::pin(async { Ok(()) })
3326 }
3327 }
3328
3329 #[tokio::test]
3330 async fn enqueue_priority_defaults_to_the_handler_priority() {
3331 let mut engine = create_test_engine();
3332 engine.register(EchoWorkflow).unwrap();
3333 engine.register(UrgentWorkflow).unwrap();
3334
3335 let echo = engine
3336 .enqueue_handler_with_options(
3337 "echo-workflow",
3338 TriggerKind::Api,
3339 json!({}),
3340 EnqueueOptions::default(),
3341 )
3342 .await
3343 .unwrap()
3344 .into_run();
3345 assert_eq!(echo.priority, 0);
3346
3347 let urgent = engine
3348 .enqueue_handler_with_options(
3349 "urgent-workflow",
3350 TriggerKind::Api,
3351 json!({}),
3352 EnqueueOptions::default(),
3353 )
3354 .await
3355 .unwrap()
3356 .into_run();
3357 assert_eq!(urgent.priority, 60);
3358 }
3359
3360 #[tokio::test]
3361 async fn enqueue_priority_explicit_value_overrides_the_handler() {
3362 let mut engine = create_test_engine();
3363 engine.register(UrgentWorkflow).unwrap();
3364
3365 let run = engine
3366 .enqueue_handler_with_options(
3367 "urgent-workflow",
3368 TriggerKind::Api,
3369 json!({}),
3370 EnqueueOptions {
3371 priority: Some(-20),
3372 ..Default::default()
3373 },
3374 )
3375 .await
3376 .unwrap()
3377 .into_run();
3378 assert_eq!(run.priority, -20);
3379
3380 let stored = engine.store().get_run(run.id).await.unwrap().unwrap();
3381 assert_eq!(stored.priority, -20);
3382 }
3383
3384 #[tokio::test]
3385 async fn enqueue_priority_out_of_range_is_rejected() {
3386 let mut engine = create_test_engine();
3387 engine.register(EchoWorkflow).unwrap();
3388
3389 for priority in [MAX_PRIORITY + 1, MIN_PRIORITY - 1] {
3390 let err = engine
3391 .enqueue_handler_with_options(
3392 "echo-workflow",
3393 TriggerKind::Api,
3394 json!({}),
3395 EnqueueOptions {
3396 priority: Some(priority),
3397 ..Default::default()
3398 },
3399 )
3400 .await
3401 .unwrap_err();
3402 assert!(matches!(err, EngineError::InvalidPriority(_)), "{err:?}");
3403 }
3404
3405 let err = engine
3407 .enqueue_handler_with_options(
3408 "not-registered",
3409 TriggerKind::Api,
3410 json!({}),
3411 EnqueueOptions {
3412 priority: Some(MAX_PRIORITY + 1),
3413 ..Default::default()
3414 },
3415 )
3416 .await
3417 .unwrap_err();
3418 assert!(matches!(err, EngineError::InvalidPriority(_)), "{err:?}");
3419
3420 let page = engine
3421 .store()
3422 .list_runs(RunFilter::default(), 1, 10)
3423 .await
3424 .unwrap();
3425 assert_eq!(page.total, 0, "no run may be created");
3426 }
3427
3428 #[tokio::test]
3429 async fn enqueue_priority_bounds_are_accepted() {
3430 let mut engine = create_test_engine();
3431 engine.register(EchoWorkflow).unwrap();
3432
3433 for priority in [MIN_PRIORITY, MAX_PRIORITY] {
3434 let run = engine
3435 .enqueue_handler_with_options(
3436 "echo-workflow",
3437 TriggerKind::Api,
3438 json!({}),
3439 EnqueueOptions {
3440 priority: Some(priority),
3441 ..Default::default()
3442 },
3443 )
3444 .await
3445 .unwrap()
3446 .into_run();
3447 assert_eq!(run.priority, priority);
3448 }
3449 }
3450
3451 struct GpuWorkflow;
3452
3453 impl WorkflowHandler for GpuWorkflow {
3454 fn name(&self) -> &str {
3455 "gpu-workflow"
3456 }
3457
3458 fn required_worker_tags(&self) -> Vec<String> {
3459 vec!["gpu".to_string()]
3460 }
3461
3462 fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3463 Box::pin(async move { Ok(()) })
3464 }
3465 }
3466
3467 #[tokio::test]
3468 async fn enqueue_merges_handler_and_request_worker_tags() {
3469 let mut engine = create_test_engine();
3470 engine.register(GpuWorkflow).unwrap();
3471
3472 let run = engine
3473 .enqueue_handler_with_options(
3474 "gpu-workflow",
3475 TriggerKind::Api,
3476 json!({}),
3477 EnqueueOptions {
3478 worker_tags: vec!["region:eu".to_string(), "gpu".to_string()],
3479 ..Default::default()
3480 },
3481 )
3482 .await
3483 .unwrap()
3484 .into_run();
3485
3486 assert_eq!(
3487 run.worker_tags,
3488 vec!["gpu".to_string(), "region:eu".to_string()]
3489 );
3490 }
3491
3492 #[tokio::test]
3493 async fn enqueue_without_worker_tags_keeps_handler_tags() {
3494 let mut engine = create_test_engine();
3495 engine.register(GpuWorkflow).unwrap();
3496 engine.register(EchoWorkflow).unwrap();
3497
3498 let gpu = engine
3499 .enqueue_handler("gpu-workflow", TriggerKind::Api, json!({}), 0)
3500 .await
3501 .unwrap();
3502 assert_eq!(gpu.worker_tags, vec!["gpu".to_string()]);
3503
3504 let echo = engine
3505 .enqueue_handler("echo-workflow", TriggerKind::Api, json!({}), 0)
3506 .await
3507 .unwrap();
3508 assert!(echo.worker_tags.is_empty());
3509 }
3510
3511 #[tokio::test]
3512 async fn enqueue_rejects_invalid_worker_tags() {
3513 let mut engine = create_test_engine();
3514 engine.register(EchoWorkflow).unwrap();
3515
3516 for worker_tags in [
3517 vec!["bad,tag".to_string()],
3518 vec![" ".to_string()],
3519 vec!["x".repeat(65)],
3520 ] {
3521 let err = engine
3522 .enqueue_handler_with_options(
3523 "echo-workflow",
3524 TriggerKind::Api,
3525 json!({}),
3526 EnqueueOptions {
3527 worker_tags,
3528 ..Default::default()
3529 },
3530 )
3531 .await
3532 .unwrap_err();
3533 assert!(matches!(err, EngineError::InvalidWorkerTag(_)), "{err:?}");
3534 }
3535
3536 let err = engine
3538 .enqueue_handler_with_options(
3539 "not-registered",
3540 TriggerKind::Api,
3541 json!({}),
3542 EnqueueOptions {
3543 worker_tags: vec!["bad,tag".to_string()],
3544 ..Default::default()
3545 },
3546 )
3547 .await
3548 .unwrap_err();
3549 assert!(matches!(err, EngineError::InvalidWorkerTag(_)), "{err:?}");
3550 }
3551
3552 #[test]
3553 fn worker_tags_are_unset_by_default() {
3554 let engine = create_test_engine();
3555 assert!(engine.worker_tags().is_none());
3556 }
3557
3558 #[test]
3559 fn set_worker_tags_stores_the_tags() {
3560 let mut engine = create_test_engine();
3561 engine.set_worker_tags(vec!["gpu".to_string()]);
3562 assert_eq!(engine.worker_tags(), Some(&["gpu".to_string()][..]));
3563
3564 engine.set_worker_tags(Vec::new());
3565 assert_eq!(engine.worker_tags(), Some(&[][..]));
3566 }
3567
3568 #[tokio::test]
3569 async fn run_handler_records_handler_worker_tags() {
3570 let mut engine = create_test_engine();
3571 engine.register(GpuWorkflow).unwrap();
3572
3573 let result = engine
3574 .run_handler("gpu-workflow", TriggerKind::Manual, json!({}))
3575 .await
3576 .unwrap();
3577 assert_eq!(result.run.worker_tags, vec!["gpu".to_string()]);
3578 }
3579
3580 #[tokio::test]
3581 async fn run_handler_leaves_the_run_unattributed() {
3582 let mut engine = create_test_engine();
3583 engine.register(EchoWorkflow).unwrap();
3584
3585 let run = engine
3586 .run_handler("echo-workflow", TriggerKind::Manual, json!({}))
3587 .await
3588 .unwrap()
3589 .run;
3590
3591 assert!(run.created_by.is_none());
3592 }
3593
3594 #[tokio::test]
3595 async fn run_handler_priority_comes_from_the_handler() {
3596 let mut engine = create_test_engine();
3597 engine.register(UrgentWorkflow).unwrap();
3598
3599 let run = engine
3600 .run_handler("urgent-workflow", TriggerKind::Manual, json!({}))
3601 .await
3602 .unwrap()
3603 .run;
3604
3605 assert_eq!(run.priority, 60);
3606 }
3607
3608 #[tokio::test]
3609 async fn engine_register_boxed() {
3610 let mut engine = create_test_engine();
3611 let handler: Box<dyn WorkflowHandler> = Box::new(EchoWorkflow);
3612 let result = engine.register_boxed(handler);
3613 assert!(result.is_ok());
3614 assert_eq!(engine.handler_names().len(), 1);
3615 }
3616
3617 #[tokio::test]
3618 async fn engine_store_and_provider_accessors() {
3619 let store = Arc::new(InMemoryStore::new());
3620 let inner = ClaudeCodeProvider::new();
3621 let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
3622 inner,
3623 "/tmp/ironflow-fixtures",
3624 ));
3625 let engine = Engine::new(store.clone(), provider.clone());
3626
3627 let _ = engine.store();
3629 let _ = engine.provider();
3630 }
3631
3632 use crate::operation::{Operation, OperationContext};
3637 use async_trait::async_trait;
3638 use ironflow_core::error::OperationError;
3639 use ironflow_store::models::StepKind;
3640
3641 struct FakeGitlabOp {
3642 project_id: u64,
3643 title: String,
3644 }
3645
3646 #[async_trait]
3647 impl Operation for FakeGitlabOp {
3648 fn kind(&self) -> &str {
3649 "gitlab"
3650 }
3651
3652 async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
3653 Ok(json!({
3654 "issue_id": 42,
3655 "project_id": self.project_id,
3656 "title": self.title,
3657 }))
3658 }
3659
3660 fn input(&self) -> Option<Value> {
3661 Some(json!({
3662 "project_id": self.project_id,
3663 "title": self.title,
3664 }))
3665 }
3666 }
3667
3668 struct FailingOp;
3669
3670 #[async_trait]
3671 impl Operation for FailingOp {
3672 fn kind(&self) -> &str {
3673 "broken-service"
3674 }
3675
3676 async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
3677 Err(OperationError::Http {
3678 status: None,
3679 message: "service unavailable".to_string(),
3680 })
3681 }
3682 }
3683
3684 struct OperationWorkflow;
3685
3686 impl WorkflowHandler for OperationWorkflow {
3687 fn name(&self) -> &str {
3688 "operation-workflow"
3689 }
3690
3691 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3692 Box::pin(async move {
3693 let op = FakeGitlabOp {
3694 project_id: 123,
3695 title: "Bug report".to_string(),
3696 };
3697 ctx.operation("create-issue", &op).await?;
3698 Ok(())
3699 })
3700 }
3701 }
3702
3703 struct FailingOperationWorkflow;
3704
3705 impl WorkflowHandler for FailingOperationWorkflow {
3706 fn name(&self) -> &str {
3707 "failing-operation-workflow"
3708 }
3709
3710 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3711 Box::pin(async move {
3712 ctx.operation("broken-call", &FailingOp).await?;
3713 Ok(())
3714 })
3715 }
3716 }
3717
3718 struct MixedWorkflow;
3719
3720 impl WorkflowHandler for MixedWorkflow {
3721 fn name(&self) -> &str {
3722 "mixed-workflow"
3723 }
3724
3725 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3726 Box::pin(async move {
3727 ctx.shell("build", ShellConfig::new("echo built")).await?;
3728 let op = FakeGitlabOp {
3729 project_id: 456,
3730 title: "Deploy done".to_string(),
3731 };
3732 let result = ctx.operation("notify-gitlab", &op).await?;
3733 assert_eq!(result.output["issue_id"], 42);
3734 Ok(())
3735 })
3736 }
3737 }
3738
3739 #[tokio::test]
3740 async fn operation_step_happy_path() {
3741 let mut engine = create_test_engine();
3742 engine.register(OperationWorkflow).unwrap();
3743
3744 let run = engine
3745 .run_handler("operation-workflow", TriggerKind::Manual, json!({}))
3746 .await
3747 .unwrap()
3748 .run;
3749
3750 assert_eq!(run.status.state, RunStatus::Completed);
3751
3752 let steps = engine.store().list_steps(run.id).await.unwrap();
3753
3754 assert_eq!(steps.len(), 1);
3755 assert_eq!(steps[0].name, "create-issue");
3756 assert_eq!(steps[0].kind, StepKind::Custom("gitlab".to_string()));
3757 assert_eq!(
3758 steps[0].status.state,
3759 ironflow_store::models::StepStatus::Completed
3760 );
3761
3762 let output = steps[0].output.as_ref().unwrap();
3763 assert_eq!(output["issue_id"], 42);
3764 assert_eq!(output["project_id"], 123);
3765
3766 let input = steps[0].input.as_ref().unwrap();
3767 assert_eq!(input["project_id"], 123);
3768 assert_eq!(input["title"], "Bug report");
3769 }
3770
3771 #[tokio::test]
3772 async fn operation_step_failure_marks_run_failed() {
3773 let mut engine = create_test_engine();
3774 engine.register(FailingOperationWorkflow).unwrap();
3775
3776 let result = engine
3777 .run_handler("failing-operation-workflow", TriggerKind::Manual, json!({}))
3778 .await;
3779
3780 assert!(result.is_err());
3781 }
3782
3783 #[tokio::test]
3784 async fn operation_mixed_with_shell_steps() {
3785 let mut engine = create_test_engine();
3786 engine.register(MixedWorkflow).unwrap();
3787
3788 let run = engine
3789 .run_handler("mixed-workflow", TriggerKind::Manual, json!({}))
3790 .await
3791 .unwrap()
3792 .run;
3793
3794 assert_eq!(run.status.state, RunStatus::Completed);
3795
3796 let steps = engine.store().list_steps(run.id).await.unwrap();
3797
3798 assert_eq!(steps.len(), 2);
3799 assert_eq!(steps[0].kind, StepKind::Shell);
3800 assert_eq!(steps[1].kind, StepKind::Custom("gitlab".to_string()));
3801 assert_eq!(steps[0].position, 0);
3802 assert_eq!(steps[1].position, 1);
3803 }
3804
3805 use crate::config::ApprovalConfig;
3810
3811 struct SingleApprovalWorkflow;
3812
3813 impl WorkflowHandler for SingleApprovalWorkflow {
3814 fn name(&self) -> &str {
3815 "single-approval"
3816 }
3817
3818 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3819 Box::pin(async move {
3820 ctx.shell("build", ShellConfig::new("echo built")).await?;
3821 ctx.approval("gate", ApprovalConfig::new("OK?")).await?;
3822 ctx.shell("deploy", ShellConfig::new("echo deployed"))
3823 .await?;
3824 Ok(())
3825 })
3826 }
3827 }
3828
3829 struct DoubleApprovalWorkflow;
3830
3831 impl WorkflowHandler for DoubleApprovalWorkflow {
3832 fn name(&self) -> &str {
3833 "double-approval"
3834 }
3835
3836 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3837 Box::pin(async move {
3838 ctx.shell("build", ShellConfig::new("echo built")).await?;
3839 ctx.approval("staging-gate", ApprovalConfig::new("Deploy staging?"))
3840 .await?;
3841 ctx.shell("deploy-staging", ShellConfig::new("echo staging"))
3842 .await?;
3843 ctx.approval("prod-gate", ApprovalConfig::new("Deploy prod?"))
3844 .await?;
3845 ctx.shell("deploy-prod", ShellConfig::new("echo prod"))
3846 .await?;
3847 Ok(())
3848 })
3849 }
3850 }
3851
3852 #[tokio::test]
3853 async fn approval_pauses_run() {
3854 let mut engine = create_test_engine();
3855 engine.register(SingleApprovalWorkflow).unwrap();
3856
3857 let run = engine
3858 .run_handler("single-approval", TriggerKind::Manual, json!({}))
3859 .await
3860 .unwrap()
3861 .run;
3862
3863 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3864
3865 let steps = engine.store().list_steps(run.id).await.unwrap();
3866 assert_eq!(steps.len(), 2); assert_eq!(steps[0].kind, StepKind::Shell);
3868 assert_eq!(steps[0].status.state, StepStatus::Completed);
3869 assert_eq!(steps[1].kind, StepKind::Approval);
3870 assert_eq!(steps[1].status.state, StepStatus::AwaitingApproval);
3871 }
3872
3873 #[tokio::test]
3874 async fn approval_resume_completes_run() {
3875 let mut engine = create_test_engine();
3876 engine.register(SingleApprovalWorkflow).unwrap();
3877
3878 let run = engine
3880 .run_handler("single-approval", TriggerKind::Manual, json!({}))
3881 .await
3882 .unwrap()
3883 .run;
3884 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3885
3886 engine
3888 .store()
3889 .update_run_status(run.id, RunStatus::Running)
3890 .await
3891 .unwrap();
3892
3893 let resumed = engine.resume_run(run.id).await.unwrap().run;
3895 assert_eq!(resumed.status.state, RunStatus::Completed);
3896
3897 let steps = engine.store().list_steps(run.id).await.unwrap();
3898 assert_eq!(steps.len(), 3); assert_eq!(steps[0].name, "build");
3900 assert_eq!(steps[0].status.state, StepStatus::Completed);
3901 assert_eq!(steps[1].name, "gate");
3902 assert_eq!(steps[1].kind, StepKind::Approval);
3903 assert_eq!(steps[1].status.state, StepStatus::Completed);
3904 assert_eq!(steps[2].name, "deploy");
3905 assert_eq!(steps[2].status.state, StepStatus::Completed);
3906 }
3907
3908 #[tokio::test]
3909 async fn double_approval_two_resumes() {
3910 let mut engine = create_test_engine();
3911 engine.register(DoubleApprovalWorkflow).unwrap();
3912
3913 let run = engine
3915 .run_handler("double-approval", TriggerKind::Manual, json!({}))
3916 .await
3917 .unwrap()
3918 .run;
3919 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3920
3921 let steps = engine.store().list_steps(run.id).await.unwrap();
3922 assert_eq!(steps.len(), 2); engine
3926 .store()
3927 .update_run_status(run.id, RunStatus::Running)
3928 .await
3929 .unwrap();
3930
3931 let resumed = engine.resume_run(run.id).await.unwrap().run;
3932 assert_eq!(resumed.status.state, RunStatus::AwaitingApproval);
3933
3934 let steps = engine.store().list_steps(run.id).await.unwrap();
3935 assert_eq!(steps.len(), 4); engine
3939 .store()
3940 .update_run_status(run.id, RunStatus::Running)
3941 .await
3942 .unwrap();
3943
3944 let final_run = engine.resume_run(run.id).await.unwrap().run;
3945 assert_eq!(final_run.status.state, RunStatus::Completed);
3946
3947 let steps = engine.store().list_steps(run.id).await.unwrap();
3948 assert_eq!(steps.len(), 5);
3949 assert_eq!(steps[0].name, "build");
3950 assert_eq!(steps[1].name, "staging-gate");
3951 assert_eq!(steps[2].name, "deploy-staging");
3952 assert_eq!(steps[3].name, "prod-gate");
3953 assert_eq!(steps[4].name, "deploy-prod");
3954
3955 for step in &steps {
3956 assert_eq!(step.status.state, StepStatus::Completed);
3957 }
3958 }
3959
3960 use ironflow_store::models::{NewStep, StepUpdate, step_trace_id};
3965
3966 async fn create_step_with_status(
3967 store: &Arc<dyn Store>,
3968 run_id: Uuid,
3969 name: &str,
3970 position: u32,
3971 status: StepStatus,
3972 ) -> ironflow_store::models::Step {
3973 let step = store
3974 .create_step(NewStep {
3975 run_id,
3976 trace_id: step_trace_id(run_id, name, position),
3977 name: name.to_string(),
3978 kind: StepKind::Shell,
3979 position,
3980 input: None,
3981 is_error_handler: false,
3982 })
3983 .await
3984 .unwrap();
3985
3986 match status {
3987 StepStatus::Pending => {}
3988 StepStatus::Running => {
3989 store
3990 .update_step(
3991 step.id,
3992 StepUpdate {
3993 status: Some(StepStatus::Running),
3994 ..StepUpdate::default()
3995 },
3996 )
3997 .await
3998 .unwrap();
3999 }
4000 StepStatus::Completed => {
4001 store
4002 .update_step(
4003 step.id,
4004 StepUpdate {
4005 status: Some(StepStatus::Running),
4006 ..StepUpdate::default()
4007 },
4008 )
4009 .await
4010 .unwrap();
4011 store
4012 .update_step(
4013 step.id,
4014 StepUpdate {
4015 status: Some(StepStatus::Completed),
4016 ..StepUpdate::default()
4017 },
4018 )
4019 .await
4020 .unwrap();
4021 }
4022 StepStatus::AwaitingApproval => {
4023 store
4024 .update_step(
4025 step.id,
4026 StepUpdate {
4027 status: Some(StepStatus::Running),
4028 ..StepUpdate::default()
4029 },
4030 )
4031 .await
4032 .unwrap();
4033 store
4034 .update_step(
4035 step.id,
4036 StepUpdate {
4037 status: Some(StepStatus::AwaitingApproval),
4038 ..StepUpdate::default()
4039 },
4040 )
4041 .await
4042 .unwrap();
4043 }
4044 _ => panic!("unsupported status for test helper: {status}"),
4045 }
4046
4047 store.get_step(step.id).await.unwrap().unwrap()
4048 }
4049
4050 #[tokio::test]
4051 async fn fail_orphaned_steps_marks_running_as_failed() {
4052 let engine = create_test_engine();
4053 let run = engine
4054 .store()
4055 .create_run(NewRun {
4056 created_by: None,
4057 workflow_name: "test".to_string(),
4058 trigger: TriggerKind::Manual,
4059 payload: json!({}),
4060 max_retries: 0,
4061 handler_version: None,
4062 labels: HashMap::new(),
4063 scheduled_at: None,
4064 idempotency_key: None,
4065 concurrency_key: None,
4066 priority: 0,
4067 concurrency_limits: Vec::new(),
4068 max_cost_usd: None,
4069 worker_tags: Vec::new(),
4070 })
4071 .await
4072 .unwrap()
4073 .into_run();
4074
4075 let step = create_step_with_status(
4076 engine.store(),
4077 run.id,
4078 "running-step",
4079 0,
4080 StepStatus::Running,
4081 )
4082 .await;
4083
4084 engine
4085 .fail_orphaned_steps(run.id, "parent run timed out")
4086 .await
4087 .unwrap();
4088
4089 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
4090 assert_eq!(updated.status.state, StepStatus::Failed);
4091 assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
4092 assert!(updated.completed_at.is_some());
4093 }
4094
4095 #[tokio::test]
4096 async fn fail_orphaned_steps_marks_pending_as_skipped() {
4097 let engine = create_test_engine();
4098 let run = engine
4099 .store()
4100 .create_run(NewRun {
4101 created_by: None,
4102 workflow_name: "test".to_string(),
4103 trigger: TriggerKind::Manual,
4104 payload: json!({}),
4105 max_retries: 0,
4106 handler_version: None,
4107 labels: HashMap::new(),
4108 scheduled_at: None,
4109 idempotency_key: None,
4110 concurrency_key: None,
4111 priority: 0,
4112 concurrency_limits: Vec::new(),
4113 max_cost_usd: None,
4114 worker_tags: Vec::new(),
4115 })
4116 .await
4117 .unwrap()
4118 .into_run();
4119
4120 let step = create_step_with_status(
4121 engine.store(),
4122 run.id,
4123 "pending-step",
4124 0,
4125 StepStatus::Pending,
4126 )
4127 .await;
4128
4129 engine
4130 .fail_orphaned_steps(run.id, "parent run timed out")
4131 .await
4132 .unwrap();
4133
4134 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
4135 assert_eq!(updated.status.state, StepStatus::Skipped);
4136 assert!(updated.error.is_none());
4137 assert!(updated.completed_at.is_some());
4138 }
4139
4140 #[tokio::test]
4141 async fn fail_orphaned_steps_marks_awaiting_approval_as_failed() {
4142 let engine = create_test_engine();
4143 let run = engine
4144 .store()
4145 .create_run(NewRun {
4146 created_by: None,
4147 workflow_name: "test".to_string(),
4148 trigger: TriggerKind::Manual,
4149 payload: json!({}),
4150 max_retries: 0,
4151 handler_version: None,
4152 labels: HashMap::new(),
4153 scheduled_at: None,
4154 idempotency_key: None,
4155 concurrency_key: None,
4156 priority: 0,
4157 concurrency_limits: Vec::new(),
4158 max_cost_usd: None,
4159 worker_tags: Vec::new(),
4160 })
4161 .await
4162 .unwrap()
4163 .into_run();
4164
4165 let step = create_step_with_status(
4166 engine.store(),
4167 run.id,
4168 "approval-step",
4169 0,
4170 StepStatus::AwaitingApproval,
4171 )
4172 .await;
4173
4174 engine
4175 .fail_orphaned_steps(run.id, "parent run timed out")
4176 .await
4177 .unwrap();
4178
4179 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
4180 assert_eq!(updated.status.state, StepStatus::Failed);
4181 assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
4182 assert!(updated.completed_at.is_some());
4183 }
4184
4185 #[tokio::test]
4186 async fn fail_orphaned_steps_skips_terminal_steps() {
4187 let engine = create_test_engine();
4188 let run = engine
4189 .store()
4190 .create_run(NewRun {
4191 created_by: None,
4192 workflow_name: "test".to_string(),
4193 trigger: TriggerKind::Manual,
4194 payload: json!({}),
4195 max_retries: 0,
4196 handler_version: None,
4197 labels: HashMap::new(),
4198 scheduled_at: None,
4199 idempotency_key: None,
4200 concurrency_key: None,
4201 priority: 0,
4202 concurrency_limits: Vec::new(),
4203 max_cost_usd: None,
4204 worker_tags: Vec::new(),
4205 })
4206 .await
4207 .unwrap()
4208 .into_run();
4209
4210 let completed_step =
4211 create_step_with_status(engine.store(), run.id, "done", 0, StepStatus::Completed).await;
4212 let running_step =
4213 create_step_with_status(engine.store(), run.id, "in-flight", 1, StepStatus::Running)
4214 .await;
4215
4216 engine
4217 .fail_orphaned_steps(run.id, "parent run timed out")
4218 .await
4219 .unwrap();
4220
4221 let completed = engine
4222 .store()
4223 .get_step(completed_step.id)
4224 .await
4225 .unwrap()
4226 .unwrap();
4227 assert_eq!(completed.status.state, StepStatus::Completed);
4228
4229 let failed = engine
4230 .store()
4231 .get_step(running_step.id)
4232 .await
4233 .unwrap()
4234 .unwrap();
4235 assert_eq!(failed.status.state, StepStatus::Failed);
4236 }
4237
4238 #[tokio::test]
4239 async fn fail_orphaned_steps_mixed_states() {
4240 let engine = create_test_engine();
4241 let run = engine
4242 .store()
4243 .create_run(NewRun {
4244 created_by: None,
4245 workflow_name: "test".to_string(),
4246 trigger: TriggerKind::Manual,
4247 payload: json!({}),
4248 max_retries: 0,
4249 handler_version: None,
4250 labels: HashMap::new(),
4251 scheduled_at: None,
4252 idempotency_key: None,
4253 concurrency_key: None,
4254 priority: 0,
4255 concurrency_limits: Vec::new(),
4256 max_cost_usd: None,
4257 worker_tags: Vec::new(),
4258 })
4259 .await
4260 .unwrap()
4261 .into_run();
4262
4263 let s_completed =
4264 create_step_with_status(engine.store(), run.id, "step-1", 0, StepStatus::Completed)
4265 .await;
4266 let s_running =
4267 create_step_with_status(engine.store(), run.id, "step-2", 1, StepStatus::Running).await;
4268 let s_pending =
4269 create_step_with_status(engine.store(), run.id, "step-3", 2, StepStatus::Pending).await;
4270
4271 engine.fail_orphaned_steps(run.id, "timeout").await.unwrap();
4272
4273 let r_completed = engine
4274 .store()
4275 .get_step(s_completed.id)
4276 .await
4277 .unwrap()
4278 .unwrap();
4279 assert_eq!(r_completed.status.state, StepStatus::Completed);
4280
4281 let r_running = engine
4282 .store()
4283 .get_step(s_running.id)
4284 .await
4285 .unwrap()
4286 .unwrap();
4287 assert_eq!(r_running.status.state, StepStatus::Failed);
4288 assert_eq!(r_running.error.as_deref(), Some("timeout"));
4289
4290 let r_pending = engine
4291 .store()
4292 .get_step(s_pending.id)
4293 .await
4294 .unwrap()
4295 .unwrap();
4296 assert_eq!(r_pending.status.state, StepStatus::Skipped);
4297 assert!(r_pending.error.is_none());
4298 }
4299
4300 #[tokio::test]
4301 async fn fail_orphaned_steps_no_steps_is_noop() {
4302 let engine = create_test_engine();
4303 let run = engine
4304 .store()
4305 .create_run(NewRun {
4306 created_by: None,
4307 workflow_name: "test".to_string(),
4308 trigger: TriggerKind::Manual,
4309 payload: json!({}),
4310 max_retries: 0,
4311 handler_version: None,
4312 labels: HashMap::new(),
4313 scheduled_at: None,
4314 idempotency_key: None,
4315 concurrency_key: None,
4316 priority: 0,
4317 concurrency_limits: Vec::new(),
4318 max_cost_usd: None,
4319 worker_tags: Vec::new(),
4320 })
4321 .await
4322 .unwrap()
4323 .into_run();
4324
4325 let result = engine.fail_orphaned_steps(run.id, "timeout").await;
4326 assert!(result.is_ok());
4327 }
4328
4329 #[tokio::test]
4330 async fn fail_orphaned_steps_preserves_existing_error() {
4331 let engine = create_test_engine();
4332 let run = engine
4333 .store()
4334 .create_run(NewRun {
4335 created_by: None,
4336 workflow_name: "test".to_string(),
4337 trigger: TriggerKind::Manual,
4338 payload: json!({}),
4339 max_retries: 0,
4340 handler_version: None,
4341 labels: HashMap::new(),
4342 scheduled_at: None,
4343 idempotency_key: None,
4344 concurrency_key: None,
4345 priority: 0,
4346 concurrency_limits: Vec::new(),
4347 max_cost_usd: None,
4348 worker_tags: Vec::new(),
4349 })
4350 .await
4351 .unwrap()
4352 .into_run();
4353
4354 let step_with_error = create_step_with_status(
4355 engine.store(),
4356 run.id,
4357 "already-errored",
4358 0,
4359 StepStatus::Running,
4360 )
4361 .await;
4362
4363 engine
4364 .store()
4365 .update_step(
4366 step_with_error.id,
4367 StepUpdate {
4368 error: Some("real error from provider".to_string()),
4369 ..StepUpdate::default()
4370 },
4371 )
4372 .await
4373 .unwrap();
4374
4375 let step_no_error = create_step_with_status(
4376 engine.store(),
4377 run.id,
4378 "no-error-yet",
4379 1,
4380 StepStatus::Running,
4381 )
4382 .await;
4383
4384 engine
4385 .fail_orphaned_steps(run.id, "parent run failed")
4386 .await
4387 .unwrap();
4388
4389 let updated_with = engine
4390 .store()
4391 .get_step(step_with_error.id)
4392 .await
4393 .unwrap()
4394 .unwrap();
4395 assert_eq!(updated_with.status.state, StepStatus::Failed);
4396 assert_eq!(
4397 updated_with.error.as_deref(),
4398 Some("real error from provider"),
4399 );
4400
4401 let updated_without = engine
4402 .store()
4403 .get_step(step_no_error.id)
4404 .await
4405 .unwrap()
4406 .unwrap();
4407 assert_eq!(updated_without.status.state, StepStatus::Failed);
4408 assert_eq!(updated_without.error.as_deref(), Some("parent run failed"),);
4409 }
4410}