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))]
998 pub async fn run_handler(
999 &self,
1000 handler_name: &str,
1001 trigger: TriggerKind,
1002 payload: Value,
1003 ) -> Result<WorkflowResult, EngineError> {
1004 let handler = self
1005 .handlers
1006 .get(handler_name)
1007 .ok_or_else(|| {
1008 EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
1009 })?
1010 .clone();
1011
1012 self.check_monthly_quota(handler_name).await?;
1013
1014 let handler_version = handler.version().map(str::to_string);
1015 let max_cost_usd = self
1016 .budget
1017 .resolve_run_cap(None, handler.default_max_cost_usd());
1018 let run = self
1019 .store
1020 .create_run(NewRun {
1021 created_by: None,
1022 workflow_name: handler_name.to_string(),
1023 trigger,
1024 payload,
1025 max_retries: 0,
1026 handler_version,
1027 labels: handler.default_labels(),
1028 scheduled_at: None,
1029 idempotency_key: None,
1030 concurrency_key: None,
1031 priority: clamp_priority(handler.priority()),
1032 concurrency_limits: Vec::new(),
1033 max_cost_usd,
1034 worker_tags: normalize_worker_tags(handler.required_worker_tags()),
1035 })
1036 .await?
1037 .into_run();
1038
1039 let run_id = run.id;
1040 info!(run_id = %run_id, handler_version = run.handler_version.as_deref().unwrap_or(""), "run created");
1041
1042 self.store
1043 .update_run_status(run_id, RunStatus::Running)
1044 .await?;
1045
1046 #[cfg(feature = "prometheus")]
1047 gauge!(RUNS_ACTIVE, "workflow" => handler_name.to_string()).increment(1.0);
1048
1049 let run_start = Instant::now();
1050 let mut ctx = self.build_context_with_guard(&run, handler.as_ref());
1051
1052 let result = handler.execute(&mut ctx).await;
1053 self.finalize_run(run_id, handler_name, result, &ctx, run_start, run.labels)
1054 .await
1055 }
1056
1057 #[tracing::instrument(name = "engine.plan_handler", skip_all, fields(workflow = %handler_name))]
1099 pub async fn plan_handler(
1100 &self,
1101 handler_name: &str,
1102 payload: Value,
1103 options: PlanOptions,
1104 ) -> Result<ExecutionPlan, EngineError> {
1105 if options.max_depth == 0 {
1106 return Err(EngineError::InvalidWorkflow(
1107 "max_depth must be at least 1".to_string(),
1108 ));
1109 }
1110
1111 let handler = self
1112 .handlers
1113 .get(handler_name)
1114 .ok_or_else(|| {
1115 EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
1116 })?
1117 .clone();
1118
1119 let estimates = if options.estimate_durations {
1120 estimate_durations(&self.store, handler_name, options.sample_runs).await?
1121 } else {
1122 HashMap::new()
1123 };
1124
1125 let shared: SharedPlanRecorder = Arc::new(Mutex::new(PlanRecorder::new(
1126 handler_name.to_string(),
1127 payload,
1128 options.max_depth,
1129 estimates,
1130 )));
1131
1132 let handlers = self.handlers.clone();
1135 let resolver: crate::context::HandlerResolver =
1136 Arc::new(move |name: &str| handlers.get(name).cloned());
1137 let mut ctx = WorkflowContext::with_handler_resolver(
1138 Uuid::now_v7(),
1139 handler_name.to_string(),
1140 self.store.clone(),
1141 self.provider.clone(),
1142 resolver,
1143 );
1144 ctx.set_plan(shared.clone());
1145
1146 if let Err(err) = handler.execute(&mut ctx).await {
1147 lock_plan(&shared).fail(err.to_string());
1148 }
1149 drop(ctx);
1150
1151 let plan = match Arc::try_unwrap(shared) {
1152 Ok(mutex) => mutex
1153 .into_inner()
1154 .unwrap_or_else(|poisoned| poisoned.into_inner())
1155 .into_plan(),
1156 Err(shared) => lock_plan(&shared).snapshot(),
1157 };
1158
1159 info!(
1160 workflow = %handler_name,
1161 steps = plan.steps.len(),
1162 truncated = plan.truncated,
1163 "execution plan built"
1164 );
1165
1166 Ok(plan)
1167 }
1168
1169 #[tracing::instrument(name = "engine.enqueue_handler", skip_all, fields(workflow = %handler_name))]
1180 pub async fn enqueue_handler(
1181 &self,
1182 handler_name: &str,
1183 trigger: TriggerKind,
1184 payload: Value,
1185 max_retries: u32,
1186 ) -> Result<Run, EngineError> {
1187 self.enqueue_handler_with_options(
1188 handler_name,
1189 trigger,
1190 payload,
1191 EnqueueOptions {
1192 max_retries,
1193 ..Default::default()
1194 },
1195 )
1196 .await
1197 .map(RunCreation::into_run)
1198 }
1199
1200 #[tracing::instrument(name = "engine.enqueue_handler_with_options", skip_all, fields(workflow = %handler_name))]
1252 pub async fn enqueue_handler_with_options(
1253 &self,
1254 handler_name: &str,
1255 trigger: TriggerKind,
1256 payload: Value,
1257 options: EnqueueOptions,
1258 ) -> Result<RunCreation, EngineError> {
1259 let EnqueueOptions {
1260 max_retries,
1261 labels,
1262 scheduled_at,
1263 max_cost_usd,
1264 created_by,
1265 idempotency_key,
1266 concurrency_key,
1267 concurrency_limits,
1268 priority,
1269 worker_tags,
1270 } = options;
1271
1272 validate_concurrency_limits(&concurrency_limits)
1275 .map_err(EngineError::InvalidConcurrencyLimit)?;
1276 if let Some(priority) = priority {
1277 validate_priority(priority).map_err(EngineError::InvalidPriority)?;
1278 }
1279 validate_worker_tags(&worker_tags).map_err(EngineError::InvalidWorkerTag)?;
1280
1281 let handler = self.handlers.get(handler_name).ok_or_else(|| {
1282 EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
1283 })?;
1284
1285 self.check_monthly_quota(handler_name).await?;
1286
1287 let handler_version = handler.version().map(str::to_string);
1288 let mut merged_labels = handler.default_labels();
1289 merged_labels.extend(labels);
1290 let resolved_cap = self
1291 .budget
1292 .resolve_run_cap(max_cost_usd, handler.default_max_cost_usd());
1293 let priority = priority.unwrap_or_else(|| clamp_priority(handler.priority()));
1294 let required_tags = normalize_worker_tags(
1295 handler
1296 .required_worker_tags()
1297 .into_iter()
1298 .chain(worker_tags),
1299 );
1300
1301 let creation = self
1302 .store
1303 .create_run(NewRun {
1304 workflow_name: handler_name.to_string(),
1305 trigger,
1306 payload,
1307 max_retries,
1308 handler_version,
1309 labels: merged_labels,
1310 scheduled_at,
1311 created_by,
1312 idempotency_key,
1313 concurrency_key,
1314 priority,
1315 concurrency_limits,
1316 max_cost_usd: resolved_cap,
1317 worker_tags: required_tags,
1318 })
1319 .await?;
1320
1321 match &creation {
1322 RunCreation::Created(run) => info!(
1323 run_id = %run.id,
1324 workflow = %handler_name,
1325 max_cost_usd = ?resolved_cap,
1326 "handler run enqueued"
1327 ),
1328 RunCreation::Existing(run) => info!(
1329 run_id = %run.id,
1330 workflow = %handler_name,
1331 "idempotent replay, nothing enqueued"
1332 ),
1333 }
1334
1335 Ok(creation)
1336 }
1337
1338 #[tracing::instrument(name = "engine.execute_handler_run", skip_all, fields(run_id = %run_id))]
1358 pub async fn execute_handler_run(&self, run_id: Uuid) -> Result<WorkflowResult, EngineError> {
1359 let run = self
1360 .store
1361 .get_run(run_id)
1362 .await?
1363 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
1364
1365 if let Some(root_run_id) = chain_root(&run) {
1366 return self.resume_chain(run, root_run_id).await;
1367 }
1368
1369 let _active = self.track_execution(run_id).await;
1370
1371 let handler = self
1372 .handlers
1373 .get(&run.workflow_name)
1374 .ok_or_else(|| {
1375 EngineError::InvalidWorkflow(format!(
1376 "no handler registered: {}",
1377 run.workflow_name
1378 ))
1379 })?
1380 .clone();
1381
1382 #[cfg(feature = "prometheus")]
1383 gauge!(RUNS_ACTIVE, "workflow" => run.workflow_name.clone()).increment(1.0);
1384
1385 let run_start = Instant::now();
1386 let mut ctx = self.build_context_with_guard(&run, handler.as_ref());
1387
1388 let result = if handler.is_version_compatible(run.handler_version.as_deref()) {
1401 ctx.load_replay_steps().await?;
1402 self.release_then_execute(run_id, handler.as_ref(), &mut ctx)
1403 .await
1404 } else {
1405 Err(EngineError::HandlerVersionMismatch {
1406 run_id,
1407 workflow_name: run.workflow_name.clone(),
1408 run_version: run
1409 .handler_version
1410 .clone()
1411 .unwrap_or_else(|| "unknown".to_string()),
1412 current_version: handler
1413 .version()
1414 .map(str::to_string)
1415 .unwrap_or_else(|| "unknown".to_string()),
1416 })
1417 };
1418
1419 self.finalize_run(
1420 run_id,
1421 &run.workflow_name,
1422 result,
1423 &ctx,
1424 run_start,
1425 run.labels,
1426 )
1427 .await
1428 }
1429
1430 #[tracing::instrument(name = "engine.execute_run", skip_all, fields(run_id = %run_id))]
1438 pub async fn execute_run(&self, run_id: Uuid) -> Result<WorkflowResult, EngineError> {
1439 self.execute_handler_run(run_id).await
1440 }
1441
1442 #[tracing::instrument(name = "engine.resume_run", skip_all, fields(run_id = %run_id))]
1471 pub async fn resume_run(&self, run_id: Uuid) -> Result<WorkflowResult, EngineError> {
1472 let run = self
1473 .store
1474 .get_run(run_id)
1475 .await?
1476 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
1477
1478 if let Some(root_run_id) = chain_root(&run) {
1479 return self.resume_chain(run, root_run_id).await;
1480 }
1481
1482 self.resume_loaded_run(run).await
1483 }
1484
1485 async fn resume_chain(
1500 &self,
1501 child: Run,
1502 root_run_id: Uuid,
1503 ) -> Result<WorkflowResult, EngineError> {
1504 let child_run_id = child.id;
1505 let lease = child.worker_id.zip(child.lease_expires_at);
1506 let root = self
1507 .store
1508 .get_run(root_run_id)
1509 .await?
1510 .ok_or(EngineError::Store(StoreError::RunNotFound(root_run_id)))?;
1511
1512 match root.status.state {
1513 RunStatus::AwaitingApproval | RunStatus::Pending => {
1514 self.move_root_to_running(root_run_id, lease.as_ref())
1515 .await?;
1516 }
1517 RunStatus::Sleeping => {
1518 self.store
1519 .update_run_status(root_run_id, RunStatus::Pending)
1520 .await?;
1521 self.move_root_to_running(root_run_id, lease.as_ref())
1522 .await?;
1523 }
1524 other => {
1525 let reason = format!(
1526 "cannot resume child run {child_run_id}: root run {root_run_id} is {other}"
1527 );
1528 if let Err(err) = self
1529 .fail_or_schedule_retry(child_run_id, &reason, false, None, None)
1530 .await
1531 {
1532 error!(
1533 run_id = %child_run_id,
1534 error = %err,
1535 "failed to fail a child run whose root cannot resume"
1536 );
1537 }
1538 return Err(EngineError::InvalidWorkflow(reason));
1539 }
1540 }
1541
1542 if lease.is_some() {
1543 self.store
1544 .update_run(
1545 child_run_id,
1546 RunUpdate {
1547 lease: Some(LeaseUpdate::Release),
1548 ..RunUpdate::default()
1549 },
1550 )
1551 .await?;
1552 }
1553
1554 info!(
1555 run_id = %child_run_id,
1556 root_run_id = %root_run_id,
1557 lease_transferred = lease.is_some(),
1558 "child run resumed through its root run"
1559 );
1560
1561 let root = self
1562 .store
1563 .get_run(root_run_id)
1564 .await?
1565 .ok_or(EngineError::Store(StoreError::RunNotFound(root_run_id)))?;
1566 self.resume_loaded_run(root).await
1567 }
1568
1569 async fn move_root_to_running(
1575 &self,
1576 root_run_id: Uuid,
1577 lease: Option<&(String, DateTime<Utc>)>,
1578 ) -> Result<(), EngineError> {
1579 match lease {
1580 Some((worker_id, expires_at)) => {
1581 self.store
1582 .update_run(
1583 root_run_id,
1584 RunUpdate {
1585 status: Some(RunStatus::Running),
1586 lease: Some(LeaseUpdate::Set {
1587 worker_id: worker_id.clone(),
1588 expires_at: *expires_at,
1589 }),
1590 ..RunUpdate::default()
1591 },
1592 )
1593 .await?;
1594 }
1595 None => {
1596 self.store
1597 .update_run_status(root_run_id, RunStatus::Running)
1598 .await?;
1599 }
1600 }
1601 Ok(())
1602 }
1603
1604 async fn resume_loaded_run(&self, run: Run) -> Result<WorkflowResult, EngineError> {
1606 let run_id = run.id;
1607 let _active = self.track_execution(run_id).await;
1608 let handler = self
1609 .handlers
1610 .get(&run.workflow_name)
1611 .ok_or_else(|| {
1612 EngineError::InvalidWorkflow(format!(
1613 "no handler registered: {}",
1614 run.workflow_name
1615 ))
1616 })?
1617 .clone();
1618
1619 info!(run_id = %run_id, workflow = %run.workflow_name, "resuming run after approval");
1620
1621 let run_start = Instant::now();
1622 let mut ctx = self.build_context_with_guard(&run, handler.as_ref());
1623
1624 let result = if handler.is_version_compatible(run.handler_version.as_deref()) {
1625 ctx.load_replay_steps().await?;
1626 self.release_then_execute(run_id, handler.as_ref(), &mut ctx)
1627 .await
1628 } else {
1629 Err(EngineError::HandlerVersionMismatch {
1630 run_id,
1631 workflow_name: run.workflow_name.clone(),
1632 run_version: run
1633 .handler_version
1634 .clone()
1635 .unwrap_or_else(|| "unknown".to_string()),
1636 current_version: handler
1637 .version()
1638 .map(str::to_string)
1639 .unwrap_or_else(|| "unknown".to_string()),
1640 })
1641 };
1642
1643 self.finalize_run(
1644 run_id,
1645 &run.workflow_name,
1646 result,
1647 &ctx,
1648 run_start,
1649 run.labels,
1650 )
1651 .await
1652 }
1653
1654 pub async fn deliver_signal(
1700 self: &Arc<Self>,
1701 signal: NewSignal,
1702 ) -> Result<SignalDelivery, EngineError> {
1703 if signal.name.trim().is_empty() {
1704 return Err(EngineError::InvalidSignal(
1705 "signal name must not be empty".to_string(),
1706 ));
1707 }
1708 if signal.key.trim().is_empty() {
1709 return Err(EngineError::InvalidSignal(
1710 "signal key must not be empty".to_string(),
1711 ));
1712 }
1713
1714 let stored = match self.store.insert_signal(signal).await? {
1715 SignalInsert::Created(stored) => stored,
1716 SignalInsert::Duplicate(existing) => {
1717 info!(
1718 signal_id = %existing.id,
1719 signal = %existing.name,
1720 key = %existing.key,
1721 "duplicate signal ignored"
1722 );
1723 return Ok(SignalDelivery {
1724 signal_id: existing.id,
1725 duplicate: true,
1726 resumed: Vec::new(),
1727 rejected: Vec::new(),
1728 });
1729 }
1730 };
1731
1732 let waiters = self
1733 .store
1734 .list_signal_waiters(&stored.name, &stored.key)
1735 .await?;
1736 let mut resumed = Vec::new();
1737 let mut rejected = Vec::new();
1738
1739 for step in waiters {
1740 if let Err(error) = validate_step_payload(step.input.as_ref(), &stored.payload) {
1741 rejected.push(SignalRejected {
1742 run_id: step.run_id,
1743 step_id: step.id,
1744 error,
1745 });
1746 continue;
1747 }
1748
1749 match self
1750 .store
1751 .resolve_signal_step(step.id, received_output(&stored))
1752 .await
1753 {
1754 Ok(SignalStepResolution::Resolved {
1755 run_id,
1756 run_resumed,
1757 }) => {
1758 resumed.push(SignalResumed {
1759 run_id,
1760 step_id: step.id,
1761 });
1762 if run_resumed && self.execution_mode == ExecutionMode::Local {
1763 self.spawn_local_resume(run_id);
1764 }
1765 }
1766 Ok(SignalStepResolution::NotWaiting { .. }) => {}
1768 Err(err) => {
1769 error!(
1770 run_id = %step.run_id,
1771 step_id = %step.id,
1772 error = %err,
1773 "failed to resolve a waiting signal step"
1774 );
1775 rejected.push(SignalRejected {
1776 run_id: step.run_id,
1777 step_id: step.id,
1778 error: err.to_string(),
1779 });
1780 }
1781 }
1782 }
1783
1784 info!(
1785 signal_id = %stored.id,
1786 signal = %stored.name,
1787 key = %stored.key,
1788 resumed = resumed.len(),
1789 rejected = rejected.len(),
1790 "signal received"
1791 );
1792 self.event_publisher
1793 .publish(Event::SignalReceived(SignalReceivedEvent {
1794 signal_id: stored.id,
1795 name: stored.name.clone(),
1796 key: stored.key.clone(),
1797 resumed_runs: resumed.iter().map(|r| r.run_id).collect(),
1798 at: stored.received_at,
1799 }));
1800
1801 Ok(SignalDelivery {
1802 signal_id: stored.id,
1803 duplicate: false,
1804 resumed,
1805 rejected,
1806 })
1807 }
1808
1809 pub async fn send_signal<S: Signal>(
1847 self: &Arc<Self>,
1848 signal: &S,
1849 key: &str,
1850 idempotency_id: Option<&str>,
1851 ) -> Result<SignalDelivery, EngineError> {
1852 let payload = to_value(signal)?;
1853 self.deliver_signal(NewSignal {
1854 name: S::NAME.to_string(),
1855 key: key.to_string(),
1856 payload,
1857 idempotency_id: idempotency_id.map(str::to_string),
1858 })
1859 .await
1860 }
1861
1862 pub(crate) fn spawn_local_resume(self: &Arc<Self>, run_id: Uuid) {
1877 let engine = Arc::clone(self);
1878 spawn(async move {
1879 engine.wait_until_idle(run_id).await;
1880 let run = match engine.store.get_run(run_id).await {
1881 Ok(Some(run)) => run,
1882 Ok(None) => {
1883 error!(run_id = %run_id, "run to restart not found");
1884 return;
1885 }
1886 Err(err) => {
1887 error!(run_id = %run_id, error = %err, "failed to load a run to restart");
1888 return;
1889 }
1890 };
1891 match run.status.state {
1892 RunStatus::Pending => {
1893 if let Err(err) = engine
1894 .store
1895 .update_run_status(run_id, RunStatus::Running)
1896 .await
1897 {
1898 error!(run_id = %run_id, error = %err, "failed to restart a woken run");
1899 return;
1900 }
1901 }
1902 RunStatus::Running => {}
1903 status => {
1904 info!(
1905 run_id = %run_id,
1906 status = %status,
1907 "run no longer waiting to restart, resume skipped"
1908 );
1909 return;
1910 }
1911 }
1912 if let Err(err) = engine.resume_run(run_id).await {
1913 error!(run_id = %run_id, error = %err, "failed to resume a woken run");
1914 }
1915 });
1916 }
1917
1918 pub async fn fail_or_schedule_retry(
1969 &self,
1970 run_id: Uuid,
1971 error: &str,
1972 retryable: bool,
1973 cost_usd: Option<Decimal>,
1974 duration_ms: Option<u64>,
1975 ) -> Result<RunStatus, EngineError> {
1976 let run = self
1977 .store
1978 .get_run(run_id)
1979 .await?
1980 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
1981
1982 if run.status.state == RunStatus::Paused {
1985 info!(run_id = %run_id, error = %error, "run paused, failure not recorded");
1986 return Ok(RunStatus::Paused);
1987 }
1988
1989 let has_attempts_left = run.retry_count < run.max_retries;
1990 let update = if retryable && has_attempts_left {
1991 let backoff = backoff_for_retry(run.retry_count);
1992 let scheduled_at = Utc::now() + TimeDelta::milliseconds(backoff.as_millis() as i64);
1993
1994 info!(
1995 run_id = %run_id,
1996 workflow = %run.workflow_name,
1997 attempt = run.retry_count + 1,
1998 max_retries = run.max_retries,
1999 backoff_secs = backoff.as_secs(),
2000 scheduled_at = %scheduled_at,
2001 "run failed, scheduling retry"
2002 );
2003
2004 RunUpdate {
2005 status: Some(RunStatus::Retrying),
2006 error: Some(error.to_string()),
2007 increment_retry: true,
2008 cost_usd,
2009 duration_ms,
2010 scheduled_at: Some(scheduled_at),
2011 ..RunUpdate::default()
2012 }
2013 } else {
2014 RunUpdate {
2015 status: Some(RunStatus::Failed),
2016 error: Some(error.to_string()),
2017 cost_usd,
2018 duration_ms,
2019 completed_at: Some(Utc::now()),
2020 ..RunUpdate::default()
2021 }
2022 };
2023
2024 let status = update.status.unwrap_or(RunStatus::Failed);
2025 self.store.update_run(run_id, update).await?;
2026 self.fail_orphaned_steps(run_id, error).await?;
2027 self.cancel_descendants_of_stopped_run(run_id, error).await;
2030
2031 Ok(status)
2032 }
2033
2034 pub async fn interrupt_running_steps(&self, run_id: Uuid) -> Result<(), EngineError> {
2064 interrupt_running_steps(self.store.as_ref(), run_id).await
2065 }
2066
2067 pub async fn fail_orphaned_steps(
2081 &self,
2082 run_id: Uuid,
2083 error_message: &str,
2084 ) -> Result<(), EngineError> {
2085 let steps = self.store.list_steps(run_id).await?;
2086 let now = Utc::now();
2087
2088 for step in steps {
2089 if step.status.state.is_terminal() {
2090 continue;
2091 }
2092
2093 let (target_status, error) = match step.status.state {
2094 StepStatus::Running | StepStatus::AwaitingApproval => {
2095 let err = if step.error.is_some() {
2096 None
2097 } else {
2098 Some(error_message.to_string())
2099 };
2100 (StepStatus::Failed, err)
2101 }
2102 StepStatus::Pending => (StepStatus::Skipped, None),
2103 _ => continue,
2104 };
2105
2106 if let Err(e) = self
2107 .store
2108 .update_step(
2109 step.id,
2110 StepUpdate {
2111 status: Some(target_status),
2112 error,
2113 completed_at: Some(now),
2114 ..StepUpdate::default()
2115 },
2116 )
2117 .await
2118 {
2119 warn!(
2120 run_id = %run_id,
2121 step_id = %step.id,
2122 step_name = %step.name,
2123 error = %e,
2124 "failed to cleanup orphaned step"
2125 );
2126 } else {
2127 info!(
2128 run_id = %run_id,
2129 step_id = %step.id,
2130 step_name = %step.name,
2131 from = %step.status.state,
2132 to = %target_status,
2133 "cleaned up orphaned step"
2134 );
2135 }
2136 }
2137
2138 Ok(())
2139 }
2140
2141 async fn release_then_execute(
2146 &self,
2147 run_id: Uuid,
2148 handler: &dyn WorkflowHandler,
2149 ctx: &mut WorkflowContext,
2150 ) -> Result<(), EngineError> {
2151 match self.provider.release_run(&run_id.to_string()).await {
2152 Ok(()) => handler.execute(ctx).await,
2153 Err(e) => Err(EngineError::Operation(OperationError::Agent(e))),
2154 }
2155 }
2156
2157 async fn finalize_run(
2163 &self,
2164 run_id: Uuid,
2165 workflow_name: &str,
2166 result: Result<(), EngineError>,
2167 ctx: &WorkflowContext,
2168 run_start: Instant,
2169 run_labels: HashMap<String, String>,
2170 ) -> Result<WorkflowResult, EngineError> {
2171 let total_duration = ctx.carried_duration_ms() + run_start.elapsed().as_millis() as u64;
2174 let completed_at = Utc::now();
2175
2176 if let Some(run) = self.store.get_run(run_id).await?
2180 && run.status.state == RunStatus::Paused
2181 {
2182 let run = self
2183 .store
2184 .update_run_returning(
2185 run_id,
2186 RunUpdate {
2187 cost_usd: Some(ctx.total_cost_usd()),
2188 duration_ms: Some(total_duration),
2189 ..RunUpdate::default()
2190 },
2191 )
2192 .await?;
2193 info!(
2194 run_id = %run_id,
2195 outcome = ?result.err().map(|err| err.to_string()),
2196 "run paused, execution stopped"
2197 );
2198 return Ok(WorkflowResult {
2199 run,
2200 steps: ctx.step_results().to_vec(),
2201 });
2202 }
2203
2204 if matches!(result, Err(EngineError::RunPaused { .. })) {
2207 let run = self
2208 .store
2209 .get_run(run_id)
2210 .await?
2211 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
2212 info!(
2213 run_id = %run_id,
2214 status = %run.status.state,
2215 "run resumed after the pause, execution stopped"
2216 );
2217 return Ok(WorkflowResult {
2218 run,
2219 steps: ctx.step_results().to_vec(),
2220 });
2221 }
2222
2223 let final_status;
2224 let final_run;
2225
2226 match result {
2227 Ok(()) => {
2228 final_status = if ctx.has_allowed_failure() {
2229 RunStatus::Warning
2230 } else {
2231 RunStatus::Completed
2232 };
2233 final_run = self
2234 .store
2235 .update_run_returning(
2236 run_id,
2237 RunUpdate {
2238 status: Some(final_status),
2239 cost_usd: Some(ctx.total_cost_usd()),
2240 duration_ms: Some(total_duration),
2241 completed_at: Some(completed_at),
2242 output: ctx.output().cloned(),
2243 ..RunUpdate::default()
2244 },
2245 )
2246 .await?;
2247
2248 info!(
2249 run_id = %run_id,
2250 status = %final_status,
2251 cost_usd = %ctx.total_cost_usd(),
2252 duration_ms = total_duration,
2253 "run completed"
2254 );
2255 }
2256 Err(EngineError::ApprovalRequired {
2257 run_id: approval_run_id,
2258 step_id,
2259 ref message,
2260 }) => {
2261 final_status = RunStatus::AwaitingApproval;
2262 final_run = self
2263 .store
2264 .update_run_returning(
2265 run_id,
2266 RunUpdate {
2267 status: Some(RunStatus::AwaitingApproval),
2268 cost_usd: Some(ctx.total_cost_usd()),
2269 duration_ms: Some(total_duration),
2270 ..RunUpdate::default()
2271 },
2272 )
2273 .await?;
2274
2275 info!(
2276 run_id = %approval_run_id,
2277 step_id = %step_id,
2278 message = %message,
2279 "run awaiting approval"
2280 );
2281
2282 self.publish_approval_requested(approval_run_id, step_id, message)
2283 .await?;
2284 }
2285 Err(EngineError::ChildSuspended {
2286 run_id: child_run_id,
2287 ref cause,
2288 }) => {
2289 final_status = cause.suspension_status();
2290 final_run = self
2294 .store
2295 .update_run_returning(
2296 run_id,
2297 RunUpdate {
2298 status: Some(final_status),
2299 cost_usd: Some(ctx.total_cost_usd()),
2300 duration_ms: Some(total_duration),
2301 ..RunUpdate::default()
2302 },
2303 )
2304 .await?;
2305
2306 let leaf = cause.suspension_leaf();
2307 info!(
2308 run_id = %run_id,
2309 child_run_id = %child_run_id,
2310 status = %final_status,
2311 cause = %leaf,
2312 "run suspended with its child run"
2313 );
2314
2315 match leaf {
2316 EngineError::ApprovalRequired {
2317 run_id: approval_run_id,
2318 step_id,
2319 message,
2320 } => {
2321 self.publish_approval_requested(*approval_run_id, *step_id, message)
2322 .await?;
2323 }
2324 EngineError::SignalWaiting {
2325 run_id: wait_run_id,
2326 step_id,
2327 step_name,
2328 name,
2329 key,
2330 deadline_at,
2331 } => {
2332 self.event_publisher
2333 .publish(Event::SignalAwaited(SignalAwaitedEvent {
2334 run_id: *wait_run_id,
2335 step_id: *step_id,
2336 step_name: step_name.clone(),
2337 name: name.clone(),
2338 key: key.clone(),
2339 deadline_at: *deadline_at,
2340 at: Utc::now(),
2341 }));
2342 }
2343 _ => {}
2346 }
2347 }
2348 Err(EngineError::HumanInputRequired {
2349 run_id: input_run_id,
2350 step_id,
2351 ref message,
2352 }) => {
2353 final_status = RunStatus::AwaitingApproval;
2354 final_run = self
2355 .store
2356 .update_run_returning(
2357 run_id,
2358 RunUpdate {
2359 status: Some(RunStatus::AwaitingApproval),
2360 cost_usd: Some(ctx.total_cost_usd()),
2361 duration_ms: Some(total_duration),
2362 ..RunUpdate::default()
2363 },
2364 )
2365 .await?;
2366
2367 info!(
2369 run_id = %input_run_id,
2370 step_id = %step_id,
2371 message = %message,
2372 "run awaiting human input"
2373 );
2374 }
2375 Err(EngineError::DelaySleeping {
2376 run_id: delay_run_id,
2377 step_id,
2378 wake_at,
2379 }) => {
2380 final_status = RunStatus::Sleeping;
2381 final_run = self
2382 .store
2383 .update_run_returning(
2384 run_id,
2385 RunUpdate {
2386 status: Some(RunStatus::Sleeping),
2387 cost_usd: Some(ctx.total_cost_usd()),
2388 duration_ms: Some(total_duration),
2389 scheduled_at: Some(wake_at),
2390 ..RunUpdate::default()
2391 },
2392 )
2393 .await?;
2394
2395 info!(
2396 run_id = %delay_run_id,
2397 step_id = %step_id,
2398 wake_at = %wake_at,
2399 "run sleeping until delay elapses"
2400 );
2401 }
2402 Err(EngineError::CapacitySleeping {
2403 run_id: capacity_run_id,
2404 step_id,
2405 ref kind,
2406 wake_at,
2407 }) => {
2408 final_status = RunStatus::Sleeping;
2409 final_run = self
2410 .store
2411 .update_run_returning(
2412 run_id,
2413 RunUpdate {
2414 status: Some(RunStatus::Sleeping),
2415 cost_usd: Some(ctx.total_cost_usd()),
2416 duration_ms: Some(total_duration),
2417 scheduled_at: Some(wake_at),
2418 capacity_wait_kind: Some(ProviderKind::new(kind.as_str())),
2419 ..RunUpdate::default()
2420 },
2421 )
2422 .await?;
2423
2424 info!(
2425 run_id = %capacity_run_id,
2426 step_id = %step_id,
2427 kind = %kind,
2428 wake_at = %wake_at,
2429 "run sleeping until provider capacity returns"
2430 );
2431 }
2432 Err(EngineError::SignalWaiting {
2433 run_id: wait_run_id,
2434 step_id,
2435 ref step_name,
2436 ref name,
2437 ref key,
2438 deadline_at,
2439 }) => {
2440 final_status = RunStatus::Sleeping;
2441 let waiting = self
2445 .store
2446 .suspend_run_on_signal(run_id, step_id, deadline_at)
2447 .await?;
2448 final_run = self
2449 .store
2450 .update_run_returning(
2451 run_id,
2452 RunUpdate {
2453 cost_usd: Some(ctx.total_cost_usd()),
2454 duration_ms: Some(total_duration),
2455 ..RunUpdate::default()
2456 },
2457 )
2458 .await?;
2459
2460 if waiting {
2461 self.event_publisher
2462 .publish(Event::SignalAwaited(SignalAwaitedEvent {
2463 run_id: wait_run_id,
2464 step_id,
2465 step_name: step_name.clone(),
2466 name: name.clone(),
2467 key: key.clone(),
2468 deadline_at,
2469 at: Utc::now(),
2470 }));
2471 }
2472
2473 info!(
2474 run_id = %wait_run_id,
2475 step_id = %step_id,
2476 signal = %name,
2477 key = %key,
2478 deadline_at = %deadline_at,
2479 waiting,
2480 "run sleeping until a signal arrives"
2481 );
2482 }
2483 Err(err) => {
2484 let guardrail_stop = matches!(
2488 err,
2489 EngineError::RunBudgetExceeded { .. } | EngineError::WorkflowGuardRejected(_)
2490 );
2491
2492 final_status = if guardrail_stop {
2493 if let Err(store_err) = self
2494 .store
2495 .update_run(
2496 run_id,
2497 RunUpdate {
2498 status: Some(RunStatus::Cancelled),
2499 error: Some(err.to_string()),
2500 cost_usd: Some(ctx.total_cost_usd()),
2501 duration_ms: Some(total_duration),
2502 completed_at: Some(completed_at),
2503 output: ctx.output().cloned(),
2504 ..RunUpdate::default()
2505 },
2506 )
2507 .await
2508 {
2509 error!(run_id = %run_id, store_error = %store_err, "failed to persist run cancellation");
2510 }
2511 if let Err(cleanup_err) = self
2512 .fail_orphaned_steps(run_id, "run stopped: guardrail limit reached")
2513 .await
2514 {
2515 error!(run_id = %run_id, store_error = %cleanup_err, "failed to cleanup orphaned steps");
2516 }
2517 RunStatus::Cancelled
2518 } else {
2519 if let Some(output) = ctx.output()
2522 && let Err(store_err) = self
2523 .store
2524 .update_run(
2525 run_id,
2526 RunUpdate {
2527 output: Some(output.clone()),
2528 ..RunUpdate::default()
2529 },
2530 )
2531 .await
2532 {
2533 error!(run_id = %run_id, store_error = %store_err, "failed to persist run output");
2534 }
2535 self.fail_or_schedule_retry(
2536 run_id,
2537 &err.to_string(),
2538 is_run_retryable(&err),
2539 Some(ctx.total_cost_usd()),
2540 Some(total_duration),
2541 )
2542 .await
2543 .unwrap_or_else(|store_err| {
2544 error!(run_id = %run_id, store_error = %store_err, "failed to persist run failure");
2545 RunStatus::Failed
2546 })
2547 };
2548
2549 if matches!(err, EngineError::RunBudgetExceeded { .. }) {
2550 self.on_run_budget_exceeded(workflow_name, run_id, &err);
2551 }
2552
2553 error!(run_id = %run_id, status = %final_status, error = %err, "run stopped");
2554
2555 self.publish_run_status_changed(
2556 workflow_name,
2557 run_id,
2558 final_status,
2559 Some(err.to_string()),
2560 ctx,
2561 total_duration,
2562 run_labels,
2563 );
2564
2565 #[cfg(feature = "prometheus")]
2566 self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
2567
2568 return Err(err);
2569 }
2570 }
2571
2572 self.publish_run_status_changed(
2573 workflow_name,
2574 run_id,
2575 final_status,
2576 None,
2577 ctx,
2578 total_duration,
2579 run_labels,
2580 );
2581
2582 #[cfg(feature = "prometheus")]
2583 self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
2584
2585 Ok(WorkflowResult {
2586 run: final_run,
2587 steps: ctx.step_results().to_vec(),
2588 })
2589 }
2590
2591 async fn publish_approval_requested(
2594 &self,
2595 run_id: Uuid,
2596 step_id: Uuid,
2597 message: &str,
2598 ) -> Result<(), EngineError> {
2599 let requirement = self
2600 .store
2601 .get_step(step_id)
2602 .await?
2603 .and_then(|s| s.approval_requirement);
2604 self.event_publisher
2605 .publish(Event::ApprovalRequested(ApprovalRequestedEvent {
2606 run_id,
2607 step_id,
2608 message: message.to_string(),
2609 requirement,
2610 at: Utc::now(),
2611 }));
2612 Ok(())
2613 }
2614
2615 pub async fn fail_ancestors(&self, run_id: Uuid, reason: &str) -> Result<(), EngineError> {
2648 let mut current = self
2649 .store
2650 .get_run(run_id)
2651 .await?
2652 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
2653 let mut visited = HashSet::from([run_id]);
2655
2656 while let Some(parent_id) = chain_parent(¤t) {
2657 if !visited.insert(parent_id) {
2658 break;
2659 }
2660 let status = self
2661 .fail_or_schedule_retry(parent_id, reason, false, None, None)
2662 .await?;
2663 if status == RunStatus::Paused {
2666 self.requeue_paused_root(¤t).await?;
2667 break;
2668 }
2669 info!(
2670 run_id = %run_id,
2671 ancestor_run_id = %parent_id,
2672 status = %status,
2673 "ancestor run failed with its child"
2674 );
2675 current = self
2676 .store
2677 .get_run(parent_id)
2678 .await?
2679 .ok_or(EngineError::Store(StoreError::RunNotFound(parent_id)))?;
2680 }
2681
2682 Ok(())
2683 }
2684
2685 #[cfg(feature = "prometheus")]
2687 fn emit_run_metrics(
2688 &self,
2689 workflow_name: &str,
2690 status: RunStatus,
2691 duration_ms: u64,
2692 ctx: &WorkflowContext,
2693 ) {
2694 let status_str = status.to_string();
2695 let wf = workflow_name.to_string();
2696
2697 counter!(RUNS_TOTAL, "workflow" => wf.clone(), "status" => status_str.clone()).increment(1);
2698 histogram!(RUN_DURATION_SECONDS, "workflow" => wf.clone(), "status" => status_str)
2699 .record(duration_ms as f64 / 1000.0);
2700 histogram!(RUN_COST_USD, "workflow" => wf.clone()).record(
2701 ctx.total_cost_usd()
2702 .to_string()
2703 .parse::<f64>()
2704 .unwrap_or(0.0),
2705 );
2706 gauge!(RUNS_ACTIVE, "workflow" => wf).decrement(1.0);
2707 }
2708
2709 fn on_run_budget_exceeded(&self, workflow_name: &str, run_id: Uuid, err: &EngineError) {
2715 let EngineError::RunBudgetExceeded {
2716 limit_usd,
2717 spent_usd,
2718 step_budget_usd,
2719 ..
2720 } = err
2721 else {
2722 return;
2723 };
2724
2725 #[cfg(feature = "prometheus")]
2726 counter!(
2727 RUN_BUDGET_EXCEEDED_TOTAL,
2728 "workflow" => workflow_name.to_string(),
2729 "scope" => "run",
2730 )
2731 .increment(1);
2732
2733 self.event_publisher
2734 .publish(Event::RunBudgetExceeded(RunBudgetExceededEvent {
2735 run_id,
2736 workflow_name: workflow_name.to_string(),
2737 limit_usd: *limit_usd,
2738 spent_usd: *spent_usd,
2739 step_budget_usd: *step_budget_usd,
2740 at: Utc::now(),
2741 }));
2742 }
2743
2744 #[allow(clippy::too_many_arguments)]
2749 fn publish_run_status_changed(
2750 &self,
2751 workflow_name: &str,
2752 run_id: Uuid,
2753 to: RunStatus,
2754 error: Option<String>,
2755 ctx: &WorkflowContext,
2756 duration_ms: u64,
2757 labels: HashMap<String, String>,
2758 ) {
2759 let now = Utc::now();
2760 let cost_usd = ctx.total_cost_usd();
2761 let wf = workflow_name.to_string();
2762
2763 self.event_publisher
2764 .publish(Event::RunStatusChanged(RunStatusChangedEvent {
2765 run_id,
2766 workflow_name: wf.clone(),
2767 from: RunStatus::Running,
2768 to,
2769 error: error.clone(),
2770 cost_usd,
2771 duration_ms,
2772 labels: labels.clone(),
2773 at: now,
2774 }));
2775
2776 if to == RunStatus::Failed {
2777 self.event_publisher
2778 .publish(Event::RunFailed(RunFailedEvent {
2779 run_id,
2780 workflow_name: wf,
2781 error,
2782 cost_usd,
2783 duration_ms,
2784 labels,
2785 at: now,
2786 }));
2787 }
2788 }
2789}
2790
2791impl fmt::Debug for Engine {
2792 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
2793 f.debug_struct("Engine")
2794 .field("handlers", &self.handlers.keys().collect::<Vec<_>>())
2795 .finish_non_exhaustive()
2796 }
2797}
2798
2799#[cfg(test)]
2800mod tests {
2801 use super::*;
2802 use crate::config::ShellConfig;
2803 use crate::handler::{HandlerFuture, WorkflowHandler};
2804 use ironflow_core::providers::claude::ClaudeCodeProvider;
2805 use ironflow_core::providers::record_replay::RecordReplayProvider;
2806 use ironflow_store::memory::InMemoryStore;
2807 use ironflow_store::models::{MAX_PRIORITY, MIN_PRIORITY, StepStatus};
2808 use serde_json::json;
2809
2810 struct EchoWorkflow;
2812
2813 impl WorkflowHandler for EchoWorkflow {
2814 fn name(&self) -> &str {
2815 "echo-workflow"
2816 }
2817
2818 fn describe(&self) -> WorkflowInfo {
2819 WorkflowInfo {
2820 description: "A simple workflow that echoes hello".to_string(),
2821 source_code: None,
2822 sub_workflows: Vec::new(),
2823 category: None,
2824 version: self.version().map(str::to_string),
2825 compatible_versions: Vec::new(),
2826 input_schema: None,
2827 default_labels: HashMap::new(),
2828 schedule: self.schedule().cloned(),
2829 default_max_cost_usd: self.default_max_cost_usd(),
2830 priority: self.priority(),
2831 }
2832 }
2833
2834 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2835 Box::pin(async move {
2836 ctx.shell("greet", ShellConfig::new("echo hello")).await?;
2837 Ok(())
2838 })
2839 }
2840 }
2841
2842 struct FailingWorkflow;
2844
2845 impl WorkflowHandler for FailingWorkflow {
2846 fn name(&self) -> &str {
2847 "failing-workflow"
2848 }
2849
2850 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2851 Box::pin(async move {
2852 ctx.shell("fail", ShellConfig::new("exit 1")).await?;
2853 Ok(())
2854 })
2855 }
2856 }
2857
2858 fn create_test_engine() -> Engine {
2859 let store = Arc::new(InMemoryStore::new());
2860 let inner = ClaudeCodeProvider::new();
2861 let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
2862 inner,
2863 "/tmp/ironflow-fixtures",
2864 ));
2865 Engine::new(store, provider)
2866 }
2867
2868 #[test]
2869 fn engine_new_creates_instance() {
2870 let engine = create_test_engine();
2871 assert_eq!(engine.handler_names().len(), 0);
2872 }
2873
2874 #[test]
2875 fn execution_mode_defaults_to_local() {
2876 let engine = create_test_engine();
2877 assert_eq!(engine.execution_mode(), ExecutionMode::Local);
2878 }
2879
2880 #[test]
2881 fn with_execution_mode_overrides_the_default() {
2882 let engine = create_test_engine().with_execution_mode(ExecutionMode::Workers);
2883 assert_eq!(engine.execution_mode(), ExecutionMode::Workers);
2884 }
2885
2886 #[test]
2887 fn engine_register_handler() {
2888 let mut engine = create_test_engine();
2889 let result = engine.register(EchoWorkflow);
2890 assert!(result.is_ok());
2891 assert_eq!(engine.handler_names().len(), 1);
2892 assert!(engine.handler_names().contains(&"echo-workflow"));
2893 }
2894
2895 #[test]
2896 fn engine_register_duplicate_returns_error() {
2897 let mut engine = create_test_engine();
2898 engine.register(EchoWorkflow).unwrap();
2899 let result = engine.register(EchoWorkflow);
2900 assert!(result.is_err());
2901 }
2902
2903 #[test]
2904 fn engine_get_handler_found() {
2905 let mut engine = create_test_engine();
2906 engine.register(EchoWorkflow).unwrap();
2907 let handler = engine.get_handler("echo-workflow");
2908 assert!(handler.is_some());
2909 }
2910
2911 #[test]
2912 fn engine_get_handler_not_found() {
2913 let engine = create_test_engine();
2914 let handler = engine.get_handler("nonexistent");
2915 assert!(handler.is_none());
2916 }
2917
2918 #[test]
2919 fn engine_handler_names_lists_all() {
2920 let mut engine = create_test_engine();
2921 engine.register(EchoWorkflow).unwrap();
2922 engine.register(FailingWorkflow).unwrap();
2923 let names = engine.handler_names();
2924 assert_eq!(names.len(), 2);
2925 assert!(names.contains(&"echo-workflow"));
2926 assert!(names.contains(&"failing-workflow"));
2927 }
2928
2929 #[test]
2930 fn engine_handler_info_returns_description() {
2931 let mut engine = create_test_engine();
2932 engine.register(EchoWorkflow).unwrap();
2933 let info = engine.handler_info("echo-workflow");
2934 assert!(info.is_some());
2935 let info = info.unwrap();
2936 assert_eq!(info.description, "A simple workflow that echoes hello");
2937 }
2938
2939 struct CategorizedWorkflow;
2940
2941 impl WorkflowHandler for CategorizedWorkflow {
2942 fn name(&self) -> &str {
2943 "categorized"
2944 }
2945 fn category(&self) -> Option<&str> {
2946 Some("data/etl")
2947 }
2948 fn execute<'a>(
2949 &'a self,
2950 _ctx: &'a mut WorkflowContext,
2951 ) -> crate::handler::HandlerFuture<'a> {
2952 Box::pin(async move { Ok(()) })
2953 }
2954 }
2955
2956 #[test]
2957 fn engine_default_describe_propagates_category() {
2958 let mut engine = create_test_engine();
2959 engine.register(CategorizedWorkflow).unwrap();
2960 let info = engine.handler_info("categorized").unwrap();
2961 assert_eq!(info.category.as_deref(), Some("data/etl"));
2962 }
2963
2964 #[test]
2965 fn engine_default_describe_without_category() {
2966 let mut engine = create_test_engine();
2967 engine.register(EchoWorkflow).unwrap();
2968 let info = engine.handler_info("echo-workflow").unwrap();
2969 assert!(info.category.is_none());
2970 }
2971
2972 struct ScheduledWorkflow {
2977 schedule: CronSchedule,
2978 }
2979
2980 impl ScheduledWorkflow {
2981 fn new() -> Self {
2982 Self {
2983 schedule: CronSchedule::new("0 0 * * * *").unwrap(),
2984 }
2985 }
2986 }
2987
2988 impl WorkflowHandler for ScheduledWorkflow {
2989 fn name(&self) -> &str {
2990 "scheduled"
2991 }
2992 fn schedule(&self) -> Option<&CronSchedule> {
2993 Some(&self.schedule)
2994 }
2995 fn execute<'a>(
2996 &'a self,
2997 _ctx: &'a mut WorkflowContext,
2998 ) -> crate::handler::HandlerFuture<'a> {
2999 Box::pin(async move { Ok(()) })
3000 }
3001 }
3002
3003 #[test]
3004 fn engine_default_describe_propagates_schedule() {
3005 let mut engine = create_test_engine();
3006 engine.register(ScheduledWorkflow::new()).unwrap();
3007 let info = engine.handler_info("scheduled").unwrap();
3008 assert_eq!(
3009 info.schedule.as_ref().map(|s| s.as_str()),
3010 Some("0 0 * * * *")
3011 );
3012 }
3013
3014 #[test]
3015 fn engine_default_describe_without_schedule() {
3016 let mut engine = create_test_engine();
3017 engine.register(EchoWorkflow).unwrap();
3018 let info = engine.handler_info("echo-workflow").unwrap();
3019 assert!(info.schedule.is_none());
3020 }
3021
3022 #[test]
3023 fn scheduled_handlers_returns_only_scheduled() {
3024 let mut engine = create_test_engine();
3025 engine.register(EchoWorkflow).unwrap();
3026 engine.register(ScheduledWorkflow::new()).unwrap();
3027 engine.register(FailingWorkflow).unwrap();
3028
3029 let scheduled = engine.scheduled_handlers();
3030 assert_eq!(scheduled.len(), 1);
3031 assert_eq!(scheduled[0].0, "scheduled");
3032 assert_eq!(scheduled[0].1.as_str(), "0 0 * * * *");
3033 }
3034
3035 #[test]
3036 fn scheduled_handlers_empty_when_none_scheduled() {
3037 let mut engine = create_test_engine();
3038 engine.register(EchoWorkflow).unwrap();
3039 engine.register(FailingWorkflow).unwrap();
3040
3041 let scheduled = engine.scheduled_handlers();
3042 assert!(scheduled.is_empty());
3043 }
3044
3045 struct BadCategoryWorkflow(&'static str);
3046
3047 impl WorkflowHandler for BadCategoryWorkflow {
3048 fn name(&self) -> &str {
3049 "bad-category"
3050 }
3051 fn category(&self) -> Option<&str> {
3052 Some(self.0)
3053 }
3054 fn execute<'a>(
3055 &'a self,
3056 _ctx: &'a mut WorkflowContext,
3057 ) -> crate::handler::HandlerFuture<'a> {
3058 Box::pin(async move { Ok(()) })
3059 }
3060 }
3061
3062 #[test]
3063 fn engine_register_rejects_empty_category() {
3064 let mut engine = create_test_engine();
3065 let err = engine.register(BadCategoryWorkflow("")).unwrap_err();
3066 match err {
3067 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty category")),
3068 other => panic!("expected InvalidWorkflow, got {other:?}"),
3069 }
3070 }
3071
3072 #[test]
3073 fn engine_register_rejects_leading_slash_category() {
3074 let mut engine = create_test_engine();
3075 let err = engine
3076 .register(BadCategoryWorkflow("/data/etl"))
3077 .unwrap_err();
3078 match err {
3079 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("leading '/'")),
3080 other => panic!("expected InvalidWorkflow, got {other:?}"),
3081 }
3082 }
3083
3084 #[test]
3085 fn engine_register_rejects_trailing_slash_category() {
3086 let mut engine = create_test_engine();
3087 let err = engine
3088 .register(BadCategoryWorkflow("data/etl/"))
3089 .unwrap_err();
3090 match err {
3091 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("trailing '/'")),
3092 other => panic!("expected InvalidWorkflow, got {other:?}"),
3093 }
3094 }
3095
3096 #[test]
3097 fn engine_register_rejects_double_slash_category() {
3098 let mut engine = create_test_engine();
3099 let err = engine
3100 .register(BadCategoryWorkflow("data//etl"))
3101 .unwrap_err();
3102 match err {
3103 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty segment")),
3104 other => panic!("expected InvalidWorkflow, got {other:?}"),
3105 }
3106 }
3107
3108 #[test]
3109 fn engine_register_rejects_whitespace_only_segment_category() {
3110 let mut engine = create_test_engine();
3111 let err = engine
3112 .register(BadCategoryWorkflow("data/ /etl"))
3113 .unwrap_err();
3114 match err {
3115 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("whitespace-only segment")),
3116 other => panic!("expected InvalidWorkflow, got {other:?}"),
3117 }
3118 }
3119
3120 #[test]
3121 fn engine_register_accepts_valid_nested_category() {
3122 let mut engine = create_test_engine();
3123 assert!(engine.register(CategorizedWorkflow).is_ok());
3124 }
3125
3126 #[tokio::test]
3127 async fn engine_unknown_workflow_returns_error() {
3128 let engine = create_test_engine();
3129 let result = engine
3130 .run_handler("unknown", TriggerKind::Manual, json!({}))
3131 .await;
3132 assert!(result.is_err());
3133 match result {
3134 Err(EngineError::InvalidWorkflow(msg)) => {
3135 assert!(msg.contains("no handler registered"));
3136 }
3137 _ => panic!("expected InvalidWorkflow error"),
3138 }
3139 }
3140
3141 #[tokio::test]
3142 async fn engine_enqueue_handler_creates_pending_run() {
3143 let mut engine = create_test_engine();
3144 engine.register(EchoWorkflow).unwrap();
3145
3146 let run = engine
3147 .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
3148 .await
3149 .unwrap();
3150 assert_eq!(run.status.state, RunStatus::Pending);
3151 assert_eq!(run.workflow_name, "echo-workflow");
3152 }
3153
3154 #[tokio::test]
3155 async fn enqueue_handler_leaves_the_run_unattributed() {
3156 let mut engine = create_test_engine();
3157 engine.register(EchoWorkflow).unwrap();
3158
3159 let run = engine
3160 .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
3161 .await
3162 .unwrap();
3163
3164 assert!(run.created_by.is_none());
3165 }
3166
3167 #[tokio::test]
3168 async fn enqueue_handler_with_options_records_the_author() {
3169 let mut engine = create_test_engine();
3170 engine.register(EchoWorkflow).unwrap();
3171 let actor = RunActor::User {
3172 user_id: Uuid::now_v7(),
3173 };
3174
3175 let run = engine
3176 .enqueue_handler_with_options(
3177 "echo-workflow",
3178 TriggerKind::Api,
3179 json!({}),
3180 EnqueueOptions {
3181 created_by: Some(actor.clone()),
3182 ..Default::default()
3183 },
3184 )
3185 .await
3186 .unwrap()
3187 .into_run();
3188
3189 assert_eq!(run.created_by, Some(actor));
3190 }
3191
3192 #[tokio::test]
3193 async fn enqueue_handler_with_options_accepts_no_author() {
3194 let mut engine = create_test_engine();
3195 engine.register(EchoWorkflow).unwrap();
3196
3197 let run = engine
3198 .enqueue_handler_with_options(
3199 "echo-workflow",
3200 TriggerKind::Cron {
3201 schedule: "0 * * * * *".to_string(),
3202 schedule_id: None,
3203 scheduled_for: None,
3204 },
3205 json!({}),
3206 EnqueueOptions::default(),
3207 )
3208 .await
3209 .unwrap()
3210 .into_run();
3211
3212 assert!(run.created_by.is_none());
3213 }
3214
3215 #[tokio::test]
3216 async fn enqueue_handler_with_options_stores_concurrency_limits() {
3217 let mut engine = create_test_engine();
3218 engine.register(EchoWorkflow).unwrap();
3219 let limits = vec![
3220 ConcurrencyLimit::new("repo:acme", 2),
3221 ConcurrencyLimit::new("tenant:42", 5),
3222 ];
3223
3224 let run = engine
3225 .enqueue_handler_with_options(
3226 "echo-workflow",
3227 TriggerKind::Api,
3228 json!({}),
3229 EnqueueOptions {
3230 concurrency_limits: limits.clone(),
3231 ..Default::default()
3232 },
3233 )
3234 .await
3235 .unwrap()
3236 .into_run();
3237
3238 assert_eq!(run.concurrency_limits, limits);
3239 }
3240
3241 #[tokio::test]
3242 async fn enqueue_rejects_invalid_concurrency_limits() {
3243 let mut engine = create_test_engine();
3244 engine.register(EchoWorkflow).unwrap();
3245
3246 let invalid = [
3247 vec![ConcurrencyLimit::new("repo:acme", 0)],
3248 vec![ConcurrencyLimit::new("", 1)],
3249 vec![
3250 ConcurrencyLimit::new("repo:acme", 1),
3251 ConcurrencyLimit::new("repo:acme", 2),
3252 ],
3253 ];
3254 for concurrency_limits in invalid {
3255 let err = engine
3256 .enqueue_handler_with_options(
3257 "echo-workflow",
3258 TriggerKind::Api,
3259 json!({}),
3260 EnqueueOptions {
3261 concurrency_limits,
3262 ..Default::default()
3263 },
3264 )
3265 .await
3266 .unwrap_err();
3267 assert!(
3268 matches!(err, EngineError::InvalidConcurrencyLimit(_)),
3269 "{err:?}"
3270 );
3271 }
3272
3273 let err = engine
3275 .enqueue_handler_with_options(
3276 "not-registered",
3277 TriggerKind::Api,
3278 json!({}),
3279 EnqueueOptions {
3280 concurrency_limits: vec![ConcurrencyLimit::new("repo:acme", 0)],
3281 ..Default::default()
3282 },
3283 )
3284 .await
3285 .unwrap_err();
3286 assert!(
3287 matches!(err, EngineError::InvalidConcurrencyLimit(_)),
3288 "{err:?}"
3289 );
3290
3291 let page = engine
3292 .store()
3293 .list_runs(RunFilter::default(), 1, 10)
3294 .await
3295 .unwrap();
3296 assert_eq!(page.total, 0, "no run may be created");
3297 }
3298
3299 struct UrgentWorkflow;
3300
3301 impl WorkflowHandler for UrgentWorkflow {
3302 fn name(&self) -> &str {
3303 "urgent-workflow"
3304 }
3305
3306 fn priority(&self) -> i16 {
3307 60
3308 }
3309
3310 fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3311 Box::pin(async { Ok(()) })
3312 }
3313 }
3314
3315 #[tokio::test]
3316 async fn enqueue_priority_defaults_to_the_handler_priority() {
3317 let mut engine = create_test_engine();
3318 engine.register(EchoWorkflow).unwrap();
3319 engine.register(UrgentWorkflow).unwrap();
3320
3321 let echo = engine
3322 .enqueue_handler_with_options(
3323 "echo-workflow",
3324 TriggerKind::Api,
3325 json!({}),
3326 EnqueueOptions::default(),
3327 )
3328 .await
3329 .unwrap()
3330 .into_run();
3331 assert_eq!(echo.priority, 0);
3332
3333 let urgent = engine
3334 .enqueue_handler_with_options(
3335 "urgent-workflow",
3336 TriggerKind::Api,
3337 json!({}),
3338 EnqueueOptions::default(),
3339 )
3340 .await
3341 .unwrap()
3342 .into_run();
3343 assert_eq!(urgent.priority, 60);
3344 }
3345
3346 #[tokio::test]
3347 async fn enqueue_priority_explicit_value_overrides_the_handler() {
3348 let mut engine = create_test_engine();
3349 engine.register(UrgentWorkflow).unwrap();
3350
3351 let run = engine
3352 .enqueue_handler_with_options(
3353 "urgent-workflow",
3354 TriggerKind::Api,
3355 json!({}),
3356 EnqueueOptions {
3357 priority: Some(-20),
3358 ..Default::default()
3359 },
3360 )
3361 .await
3362 .unwrap()
3363 .into_run();
3364 assert_eq!(run.priority, -20);
3365
3366 let stored = engine.store().get_run(run.id).await.unwrap().unwrap();
3367 assert_eq!(stored.priority, -20);
3368 }
3369
3370 #[tokio::test]
3371 async fn enqueue_priority_out_of_range_is_rejected() {
3372 let mut engine = create_test_engine();
3373 engine.register(EchoWorkflow).unwrap();
3374
3375 for priority in [MAX_PRIORITY + 1, MIN_PRIORITY - 1] {
3376 let err = engine
3377 .enqueue_handler_with_options(
3378 "echo-workflow",
3379 TriggerKind::Api,
3380 json!({}),
3381 EnqueueOptions {
3382 priority: Some(priority),
3383 ..Default::default()
3384 },
3385 )
3386 .await
3387 .unwrap_err();
3388 assert!(matches!(err, EngineError::InvalidPriority(_)), "{err:?}");
3389 }
3390
3391 let err = engine
3393 .enqueue_handler_with_options(
3394 "not-registered",
3395 TriggerKind::Api,
3396 json!({}),
3397 EnqueueOptions {
3398 priority: Some(MAX_PRIORITY + 1),
3399 ..Default::default()
3400 },
3401 )
3402 .await
3403 .unwrap_err();
3404 assert!(matches!(err, EngineError::InvalidPriority(_)), "{err:?}");
3405
3406 let page = engine
3407 .store()
3408 .list_runs(RunFilter::default(), 1, 10)
3409 .await
3410 .unwrap();
3411 assert_eq!(page.total, 0, "no run may be created");
3412 }
3413
3414 #[tokio::test]
3415 async fn enqueue_priority_bounds_are_accepted() {
3416 let mut engine = create_test_engine();
3417 engine.register(EchoWorkflow).unwrap();
3418
3419 for priority in [MIN_PRIORITY, MAX_PRIORITY] {
3420 let run = engine
3421 .enqueue_handler_with_options(
3422 "echo-workflow",
3423 TriggerKind::Api,
3424 json!({}),
3425 EnqueueOptions {
3426 priority: Some(priority),
3427 ..Default::default()
3428 },
3429 )
3430 .await
3431 .unwrap()
3432 .into_run();
3433 assert_eq!(run.priority, priority);
3434 }
3435 }
3436
3437 struct GpuWorkflow;
3438
3439 impl WorkflowHandler for GpuWorkflow {
3440 fn name(&self) -> &str {
3441 "gpu-workflow"
3442 }
3443
3444 fn required_worker_tags(&self) -> Vec<String> {
3445 vec!["gpu".to_string()]
3446 }
3447
3448 fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3449 Box::pin(async move { Ok(()) })
3450 }
3451 }
3452
3453 #[tokio::test]
3454 async fn enqueue_merges_handler_and_request_worker_tags() {
3455 let mut engine = create_test_engine();
3456 engine.register(GpuWorkflow).unwrap();
3457
3458 let run = engine
3459 .enqueue_handler_with_options(
3460 "gpu-workflow",
3461 TriggerKind::Api,
3462 json!({}),
3463 EnqueueOptions {
3464 worker_tags: vec!["region:eu".to_string(), "gpu".to_string()],
3465 ..Default::default()
3466 },
3467 )
3468 .await
3469 .unwrap()
3470 .into_run();
3471
3472 assert_eq!(
3473 run.worker_tags,
3474 vec!["gpu".to_string(), "region:eu".to_string()]
3475 );
3476 }
3477
3478 #[tokio::test]
3479 async fn enqueue_without_worker_tags_keeps_handler_tags() {
3480 let mut engine = create_test_engine();
3481 engine.register(GpuWorkflow).unwrap();
3482 engine.register(EchoWorkflow).unwrap();
3483
3484 let gpu = engine
3485 .enqueue_handler("gpu-workflow", TriggerKind::Api, json!({}), 0)
3486 .await
3487 .unwrap();
3488 assert_eq!(gpu.worker_tags, vec!["gpu".to_string()]);
3489
3490 let echo = engine
3491 .enqueue_handler("echo-workflow", TriggerKind::Api, json!({}), 0)
3492 .await
3493 .unwrap();
3494 assert!(echo.worker_tags.is_empty());
3495 }
3496
3497 #[tokio::test]
3498 async fn enqueue_rejects_invalid_worker_tags() {
3499 let mut engine = create_test_engine();
3500 engine.register(EchoWorkflow).unwrap();
3501
3502 for worker_tags in [
3503 vec!["bad,tag".to_string()],
3504 vec![" ".to_string()],
3505 vec!["x".repeat(65)],
3506 ] {
3507 let err = engine
3508 .enqueue_handler_with_options(
3509 "echo-workflow",
3510 TriggerKind::Api,
3511 json!({}),
3512 EnqueueOptions {
3513 worker_tags,
3514 ..Default::default()
3515 },
3516 )
3517 .await
3518 .unwrap_err();
3519 assert!(matches!(err, EngineError::InvalidWorkerTag(_)), "{err:?}");
3520 }
3521
3522 let err = engine
3524 .enqueue_handler_with_options(
3525 "not-registered",
3526 TriggerKind::Api,
3527 json!({}),
3528 EnqueueOptions {
3529 worker_tags: vec!["bad,tag".to_string()],
3530 ..Default::default()
3531 },
3532 )
3533 .await
3534 .unwrap_err();
3535 assert!(matches!(err, EngineError::InvalidWorkerTag(_)), "{err:?}");
3536 }
3537
3538 #[test]
3539 fn worker_tags_are_unset_by_default() {
3540 let engine = create_test_engine();
3541 assert!(engine.worker_tags().is_none());
3542 }
3543
3544 #[test]
3545 fn set_worker_tags_stores_the_tags() {
3546 let mut engine = create_test_engine();
3547 engine.set_worker_tags(vec!["gpu".to_string()]);
3548 assert_eq!(engine.worker_tags(), Some(&["gpu".to_string()][..]));
3549
3550 engine.set_worker_tags(Vec::new());
3551 assert_eq!(engine.worker_tags(), Some(&[][..]));
3552 }
3553
3554 #[tokio::test]
3555 async fn run_handler_records_handler_worker_tags() {
3556 let mut engine = create_test_engine();
3557 engine.register(GpuWorkflow).unwrap();
3558
3559 let result = engine
3560 .run_handler("gpu-workflow", TriggerKind::Manual, json!({}))
3561 .await
3562 .unwrap();
3563 assert_eq!(result.run.worker_tags, vec!["gpu".to_string()]);
3564 }
3565
3566 #[tokio::test]
3567 async fn run_handler_leaves_the_run_unattributed() {
3568 let mut engine = create_test_engine();
3569 engine.register(EchoWorkflow).unwrap();
3570
3571 let run = engine
3572 .run_handler("echo-workflow", TriggerKind::Manual, json!({}))
3573 .await
3574 .unwrap()
3575 .run;
3576
3577 assert!(run.created_by.is_none());
3578 }
3579
3580 #[tokio::test]
3581 async fn run_handler_priority_comes_from_the_handler() {
3582 let mut engine = create_test_engine();
3583 engine.register(UrgentWorkflow).unwrap();
3584
3585 let run = engine
3586 .run_handler("urgent-workflow", TriggerKind::Manual, json!({}))
3587 .await
3588 .unwrap()
3589 .run;
3590
3591 assert_eq!(run.priority, 60);
3592 }
3593
3594 #[tokio::test]
3595 async fn engine_register_boxed() {
3596 let mut engine = create_test_engine();
3597 let handler: Box<dyn WorkflowHandler> = Box::new(EchoWorkflow);
3598 let result = engine.register_boxed(handler);
3599 assert!(result.is_ok());
3600 assert_eq!(engine.handler_names().len(), 1);
3601 }
3602
3603 #[tokio::test]
3604 async fn engine_store_and_provider_accessors() {
3605 let store = Arc::new(InMemoryStore::new());
3606 let inner = ClaudeCodeProvider::new();
3607 let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
3608 inner,
3609 "/tmp/ironflow-fixtures",
3610 ));
3611 let engine = Engine::new(store.clone(), provider.clone());
3612
3613 let _ = engine.store();
3615 let _ = engine.provider();
3616 }
3617
3618 use crate::operation::{Operation, OperationContext};
3623 use async_trait::async_trait;
3624 use ironflow_core::error::OperationError;
3625 use ironflow_store::models::StepKind;
3626
3627 struct FakeGitlabOp {
3628 project_id: u64,
3629 title: String,
3630 }
3631
3632 #[async_trait]
3633 impl Operation for FakeGitlabOp {
3634 fn kind(&self) -> &str {
3635 "gitlab"
3636 }
3637
3638 async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
3639 Ok(json!({
3640 "issue_id": 42,
3641 "project_id": self.project_id,
3642 "title": self.title,
3643 }))
3644 }
3645
3646 fn input(&self) -> Option<Value> {
3647 Some(json!({
3648 "project_id": self.project_id,
3649 "title": self.title,
3650 }))
3651 }
3652 }
3653
3654 struct FailingOp;
3655
3656 #[async_trait]
3657 impl Operation for FailingOp {
3658 fn kind(&self) -> &str {
3659 "broken-service"
3660 }
3661
3662 async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
3663 Err(OperationError::Http {
3664 status: None,
3665 message: "service unavailable".to_string(),
3666 })
3667 }
3668 }
3669
3670 struct OperationWorkflow;
3671
3672 impl WorkflowHandler for OperationWorkflow {
3673 fn name(&self) -> &str {
3674 "operation-workflow"
3675 }
3676
3677 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3678 Box::pin(async move {
3679 let op = FakeGitlabOp {
3680 project_id: 123,
3681 title: "Bug report".to_string(),
3682 };
3683 ctx.operation("create-issue", &op).await?;
3684 Ok(())
3685 })
3686 }
3687 }
3688
3689 struct FailingOperationWorkflow;
3690
3691 impl WorkflowHandler for FailingOperationWorkflow {
3692 fn name(&self) -> &str {
3693 "failing-operation-workflow"
3694 }
3695
3696 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3697 Box::pin(async move {
3698 ctx.operation("broken-call", &FailingOp).await?;
3699 Ok(())
3700 })
3701 }
3702 }
3703
3704 struct MixedWorkflow;
3705
3706 impl WorkflowHandler for MixedWorkflow {
3707 fn name(&self) -> &str {
3708 "mixed-workflow"
3709 }
3710
3711 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3712 Box::pin(async move {
3713 ctx.shell("build", ShellConfig::new("echo built")).await?;
3714 let op = FakeGitlabOp {
3715 project_id: 456,
3716 title: "Deploy done".to_string(),
3717 };
3718 let result = ctx.operation("notify-gitlab", &op).await?;
3719 assert_eq!(result.output["issue_id"], 42);
3720 Ok(())
3721 })
3722 }
3723 }
3724
3725 #[tokio::test]
3726 async fn operation_step_happy_path() {
3727 let mut engine = create_test_engine();
3728 engine.register(OperationWorkflow).unwrap();
3729
3730 let run = engine
3731 .run_handler("operation-workflow", TriggerKind::Manual, json!({}))
3732 .await
3733 .unwrap()
3734 .run;
3735
3736 assert_eq!(run.status.state, RunStatus::Completed);
3737
3738 let steps = engine.store().list_steps(run.id).await.unwrap();
3739
3740 assert_eq!(steps.len(), 1);
3741 assert_eq!(steps[0].name, "create-issue");
3742 assert_eq!(steps[0].kind, StepKind::Custom("gitlab".to_string()));
3743 assert_eq!(
3744 steps[0].status.state,
3745 ironflow_store::models::StepStatus::Completed
3746 );
3747
3748 let output = steps[0].output.as_ref().unwrap();
3749 assert_eq!(output["issue_id"], 42);
3750 assert_eq!(output["project_id"], 123);
3751
3752 let input = steps[0].input.as_ref().unwrap();
3753 assert_eq!(input["project_id"], 123);
3754 assert_eq!(input["title"], "Bug report");
3755 }
3756
3757 #[tokio::test]
3758 async fn operation_step_failure_marks_run_failed() {
3759 let mut engine = create_test_engine();
3760 engine.register(FailingOperationWorkflow).unwrap();
3761
3762 let result = engine
3763 .run_handler("failing-operation-workflow", TriggerKind::Manual, json!({}))
3764 .await;
3765
3766 assert!(result.is_err());
3767 }
3768
3769 #[tokio::test]
3770 async fn operation_mixed_with_shell_steps() {
3771 let mut engine = create_test_engine();
3772 engine.register(MixedWorkflow).unwrap();
3773
3774 let run = engine
3775 .run_handler("mixed-workflow", TriggerKind::Manual, json!({}))
3776 .await
3777 .unwrap()
3778 .run;
3779
3780 assert_eq!(run.status.state, RunStatus::Completed);
3781
3782 let steps = engine.store().list_steps(run.id).await.unwrap();
3783
3784 assert_eq!(steps.len(), 2);
3785 assert_eq!(steps[0].kind, StepKind::Shell);
3786 assert_eq!(steps[1].kind, StepKind::Custom("gitlab".to_string()));
3787 assert_eq!(steps[0].position, 0);
3788 assert_eq!(steps[1].position, 1);
3789 }
3790
3791 use crate::config::ApprovalConfig;
3796
3797 struct SingleApprovalWorkflow;
3798
3799 impl WorkflowHandler for SingleApprovalWorkflow {
3800 fn name(&self) -> &str {
3801 "single-approval"
3802 }
3803
3804 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3805 Box::pin(async move {
3806 ctx.shell("build", ShellConfig::new("echo built")).await?;
3807 ctx.approval("gate", ApprovalConfig::new("OK?")).await?;
3808 ctx.shell("deploy", ShellConfig::new("echo deployed"))
3809 .await?;
3810 Ok(())
3811 })
3812 }
3813 }
3814
3815 struct DoubleApprovalWorkflow;
3816
3817 impl WorkflowHandler for DoubleApprovalWorkflow {
3818 fn name(&self) -> &str {
3819 "double-approval"
3820 }
3821
3822 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3823 Box::pin(async move {
3824 ctx.shell("build", ShellConfig::new("echo built")).await?;
3825 ctx.approval("staging-gate", ApprovalConfig::new("Deploy staging?"))
3826 .await?;
3827 ctx.shell("deploy-staging", ShellConfig::new("echo staging"))
3828 .await?;
3829 ctx.approval("prod-gate", ApprovalConfig::new("Deploy prod?"))
3830 .await?;
3831 ctx.shell("deploy-prod", ShellConfig::new("echo prod"))
3832 .await?;
3833 Ok(())
3834 })
3835 }
3836 }
3837
3838 #[tokio::test]
3839 async fn approval_pauses_run() {
3840 let mut engine = create_test_engine();
3841 engine.register(SingleApprovalWorkflow).unwrap();
3842
3843 let run = engine
3844 .run_handler("single-approval", TriggerKind::Manual, json!({}))
3845 .await
3846 .unwrap()
3847 .run;
3848
3849 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3850
3851 let steps = engine.store().list_steps(run.id).await.unwrap();
3852 assert_eq!(steps.len(), 2); assert_eq!(steps[0].kind, StepKind::Shell);
3854 assert_eq!(steps[0].status.state, StepStatus::Completed);
3855 assert_eq!(steps[1].kind, StepKind::Approval);
3856 assert_eq!(steps[1].status.state, StepStatus::AwaitingApproval);
3857 }
3858
3859 #[tokio::test]
3860 async fn approval_resume_completes_run() {
3861 let mut engine = create_test_engine();
3862 engine.register(SingleApprovalWorkflow).unwrap();
3863
3864 let run = engine
3866 .run_handler("single-approval", TriggerKind::Manual, json!({}))
3867 .await
3868 .unwrap()
3869 .run;
3870 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3871
3872 engine
3874 .store()
3875 .update_run_status(run.id, RunStatus::Running)
3876 .await
3877 .unwrap();
3878
3879 let resumed = engine.resume_run(run.id).await.unwrap().run;
3881 assert_eq!(resumed.status.state, RunStatus::Completed);
3882
3883 let steps = engine.store().list_steps(run.id).await.unwrap();
3884 assert_eq!(steps.len(), 3); assert_eq!(steps[0].name, "build");
3886 assert_eq!(steps[0].status.state, StepStatus::Completed);
3887 assert_eq!(steps[1].name, "gate");
3888 assert_eq!(steps[1].kind, StepKind::Approval);
3889 assert_eq!(steps[1].status.state, StepStatus::Completed);
3890 assert_eq!(steps[2].name, "deploy");
3891 assert_eq!(steps[2].status.state, StepStatus::Completed);
3892 }
3893
3894 #[tokio::test]
3895 async fn double_approval_two_resumes() {
3896 let mut engine = create_test_engine();
3897 engine.register(DoubleApprovalWorkflow).unwrap();
3898
3899 let run = engine
3901 .run_handler("double-approval", TriggerKind::Manual, json!({}))
3902 .await
3903 .unwrap()
3904 .run;
3905 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3906
3907 let steps = engine.store().list_steps(run.id).await.unwrap();
3908 assert_eq!(steps.len(), 2); engine
3912 .store()
3913 .update_run_status(run.id, RunStatus::Running)
3914 .await
3915 .unwrap();
3916
3917 let resumed = engine.resume_run(run.id).await.unwrap().run;
3918 assert_eq!(resumed.status.state, RunStatus::AwaitingApproval);
3919
3920 let steps = engine.store().list_steps(run.id).await.unwrap();
3921 assert_eq!(steps.len(), 4); engine
3925 .store()
3926 .update_run_status(run.id, RunStatus::Running)
3927 .await
3928 .unwrap();
3929
3930 let final_run = engine.resume_run(run.id).await.unwrap().run;
3931 assert_eq!(final_run.status.state, RunStatus::Completed);
3932
3933 let steps = engine.store().list_steps(run.id).await.unwrap();
3934 assert_eq!(steps.len(), 5);
3935 assert_eq!(steps[0].name, "build");
3936 assert_eq!(steps[1].name, "staging-gate");
3937 assert_eq!(steps[2].name, "deploy-staging");
3938 assert_eq!(steps[3].name, "prod-gate");
3939 assert_eq!(steps[4].name, "deploy-prod");
3940
3941 for step in &steps {
3942 assert_eq!(step.status.state, StepStatus::Completed);
3943 }
3944 }
3945
3946 use ironflow_store::models::{NewStep, StepUpdate, step_trace_id};
3951
3952 async fn create_step_with_status(
3953 store: &Arc<dyn Store>,
3954 run_id: Uuid,
3955 name: &str,
3956 position: u32,
3957 status: StepStatus,
3958 ) -> ironflow_store::models::Step {
3959 let step = store
3960 .create_step(NewStep {
3961 run_id,
3962 trace_id: step_trace_id(run_id, name, position),
3963 name: name.to_string(),
3964 kind: StepKind::Shell,
3965 position,
3966 input: None,
3967 is_error_handler: false,
3968 })
3969 .await
3970 .unwrap();
3971
3972 match status {
3973 StepStatus::Pending => {}
3974 StepStatus::Running => {
3975 store
3976 .update_step(
3977 step.id,
3978 StepUpdate {
3979 status: Some(StepStatus::Running),
3980 ..StepUpdate::default()
3981 },
3982 )
3983 .await
3984 .unwrap();
3985 }
3986 StepStatus::Completed => {
3987 store
3988 .update_step(
3989 step.id,
3990 StepUpdate {
3991 status: Some(StepStatus::Running),
3992 ..StepUpdate::default()
3993 },
3994 )
3995 .await
3996 .unwrap();
3997 store
3998 .update_step(
3999 step.id,
4000 StepUpdate {
4001 status: Some(StepStatus::Completed),
4002 ..StepUpdate::default()
4003 },
4004 )
4005 .await
4006 .unwrap();
4007 }
4008 StepStatus::AwaitingApproval => {
4009 store
4010 .update_step(
4011 step.id,
4012 StepUpdate {
4013 status: Some(StepStatus::Running),
4014 ..StepUpdate::default()
4015 },
4016 )
4017 .await
4018 .unwrap();
4019 store
4020 .update_step(
4021 step.id,
4022 StepUpdate {
4023 status: Some(StepStatus::AwaitingApproval),
4024 ..StepUpdate::default()
4025 },
4026 )
4027 .await
4028 .unwrap();
4029 }
4030 _ => panic!("unsupported status for test helper: {status}"),
4031 }
4032
4033 store.get_step(step.id).await.unwrap().unwrap()
4034 }
4035
4036 #[tokio::test]
4037 async fn fail_orphaned_steps_marks_running_as_failed() {
4038 let engine = create_test_engine();
4039 let run = engine
4040 .store()
4041 .create_run(NewRun {
4042 created_by: None,
4043 workflow_name: "test".to_string(),
4044 trigger: TriggerKind::Manual,
4045 payload: json!({}),
4046 max_retries: 0,
4047 handler_version: None,
4048 labels: HashMap::new(),
4049 scheduled_at: None,
4050 idempotency_key: None,
4051 concurrency_key: None,
4052 priority: 0,
4053 concurrency_limits: Vec::new(),
4054 max_cost_usd: None,
4055 worker_tags: Vec::new(),
4056 })
4057 .await
4058 .unwrap()
4059 .into_run();
4060
4061 let step = create_step_with_status(
4062 engine.store(),
4063 run.id,
4064 "running-step",
4065 0,
4066 StepStatus::Running,
4067 )
4068 .await;
4069
4070 engine
4071 .fail_orphaned_steps(run.id, "parent run timed out")
4072 .await
4073 .unwrap();
4074
4075 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
4076 assert_eq!(updated.status.state, StepStatus::Failed);
4077 assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
4078 assert!(updated.completed_at.is_some());
4079 }
4080
4081 #[tokio::test]
4082 async fn fail_orphaned_steps_marks_pending_as_skipped() {
4083 let engine = create_test_engine();
4084 let run = engine
4085 .store()
4086 .create_run(NewRun {
4087 created_by: None,
4088 workflow_name: "test".to_string(),
4089 trigger: TriggerKind::Manual,
4090 payload: json!({}),
4091 max_retries: 0,
4092 handler_version: None,
4093 labels: HashMap::new(),
4094 scheduled_at: None,
4095 idempotency_key: None,
4096 concurrency_key: None,
4097 priority: 0,
4098 concurrency_limits: Vec::new(),
4099 max_cost_usd: None,
4100 worker_tags: Vec::new(),
4101 })
4102 .await
4103 .unwrap()
4104 .into_run();
4105
4106 let step = create_step_with_status(
4107 engine.store(),
4108 run.id,
4109 "pending-step",
4110 0,
4111 StepStatus::Pending,
4112 )
4113 .await;
4114
4115 engine
4116 .fail_orphaned_steps(run.id, "parent run timed out")
4117 .await
4118 .unwrap();
4119
4120 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
4121 assert_eq!(updated.status.state, StepStatus::Skipped);
4122 assert!(updated.error.is_none());
4123 assert!(updated.completed_at.is_some());
4124 }
4125
4126 #[tokio::test]
4127 async fn fail_orphaned_steps_marks_awaiting_approval_as_failed() {
4128 let engine = create_test_engine();
4129 let run = engine
4130 .store()
4131 .create_run(NewRun {
4132 created_by: None,
4133 workflow_name: "test".to_string(),
4134 trigger: TriggerKind::Manual,
4135 payload: json!({}),
4136 max_retries: 0,
4137 handler_version: None,
4138 labels: HashMap::new(),
4139 scheduled_at: None,
4140 idempotency_key: None,
4141 concurrency_key: None,
4142 priority: 0,
4143 concurrency_limits: Vec::new(),
4144 max_cost_usd: None,
4145 worker_tags: Vec::new(),
4146 })
4147 .await
4148 .unwrap()
4149 .into_run();
4150
4151 let step = create_step_with_status(
4152 engine.store(),
4153 run.id,
4154 "approval-step",
4155 0,
4156 StepStatus::AwaitingApproval,
4157 )
4158 .await;
4159
4160 engine
4161 .fail_orphaned_steps(run.id, "parent run timed out")
4162 .await
4163 .unwrap();
4164
4165 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
4166 assert_eq!(updated.status.state, StepStatus::Failed);
4167 assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
4168 assert!(updated.completed_at.is_some());
4169 }
4170
4171 #[tokio::test]
4172 async fn fail_orphaned_steps_skips_terminal_steps() {
4173 let engine = create_test_engine();
4174 let run = engine
4175 .store()
4176 .create_run(NewRun {
4177 created_by: None,
4178 workflow_name: "test".to_string(),
4179 trigger: TriggerKind::Manual,
4180 payload: json!({}),
4181 max_retries: 0,
4182 handler_version: None,
4183 labels: HashMap::new(),
4184 scheduled_at: None,
4185 idempotency_key: None,
4186 concurrency_key: None,
4187 priority: 0,
4188 concurrency_limits: Vec::new(),
4189 max_cost_usd: None,
4190 worker_tags: Vec::new(),
4191 })
4192 .await
4193 .unwrap()
4194 .into_run();
4195
4196 let completed_step =
4197 create_step_with_status(engine.store(), run.id, "done", 0, StepStatus::Completed).await;
4198 let running_step =
4199 create_step_with_status(engine.store(), run.id, "in-flight", 1, StepStatus::Running)
4200 .await;
4201
4202 engine
4203 .fail_orphaned_steps(run.id, "parent run timed out")
4204 .await
4205 .unwrap();
4206
4207 let completed = engine
4208 .store()
4209 .get_step(completed_step.id)
4210 .await
4211 .unwrap()
4212 .unwrap();
4213 assert_eq!(completed.status.state, StepStatus::Completed);
4214
4215 let failed = engine
4216 .store()
4217 .get_step(running_step.id)
4218 .await
4219 .unwrap()
4220 .unwrap();
4221 assert_eq!(failed.status.state, StepStatus::Failed);
4222 }
4223
4224 #[tokio::test]
4225 async fn fail_orphaned_steps_mixed_states() {
4226 let engine = create_test_engine();
4227 let run = engine
4228 .store()
4229 .create_run(NewRun {
4230 created_by: None,
4231 workflow_name: "test".to_string(),
4232 trigger: TriggerKind::Manual,
4233 payload: json!({}),
4234 max_retries: 0,
4235 handler_version: None,
4236 labels: HashMap::new(),
4237 scheduled_at: None,
4238 idempotency_key: None,
4239 concurrency_key: None,
4240 priority: 0,
4241 concurrency_limits: Vec::new(),
4242 max_cost_usd: None,
4243 worker_tags: Vec::new(),
4244 })
4245 .await
4246 .unwrap()
4247 .into_run();
4248
4249 let s_completed =
4250 create_step_with_status(engine.store(), run.id, "step-1", 0, StepStatus::Completed)
4251 .await;
4252 let s_running =
4253 create_step_with_status(engine.store(), run.id, "step-2", 1, StepStatus::Running).await;
4254 let s_pending =
4255 create_step_with_status(engine.store(), run.id, "step-3", 2, StepStatus::Pending).await;
4256
4257 engine.fail_orphaned_steps(run.id, "timeout").await.unwrap();
4258
4259 let r_completed = engine
4260 .store()
4261 .get_step(s_completed.id)
4262 .await
4263 .unwrap()
4264 .unwrap();
4265 assert_eq!(r_completed.status.state, StepStatus::Completed);
4266
4267 let r_running = engine
4268 .store()
4269 .get_step(s_running.id)
4270 .await
4271 .unwrap()
4272 .unwrap();
4273 assert_eq!(r_running.status.state, StepStatus::Failed);
4274 assert_eq!(r_running.error.as_deref(), Some("timeout"));
4275
4276 let r_pending = engine
4277 .store()
4278 .get_step(s_pending.id)
4279 .await
4280 .unwrap()
4281 .unwrap();
4282 assert_eq!(r_pending.status.state, StepStatus::Skipped);
4283 assert!(r_pending.error.is_none());
4284 }
4285
4286 #[tokio::test]
4287 async fn fail_orphaned_steps_no_steps_is_noop() {
4288 let engine = create_test_engine();
4289 let run = engine
4290 .store()
4291 .create_run(NewRun {
4292 created_by: None,
4293 workflow_name: "test".to_string(),
4294 trigger: TriggerKind::Manual,
4295 payload: json!({}),
4296 max_retries: 0,
4297 handler_version: None,
4298 labels: HashMap::new(),
4299 scheduled_at: None,
4300 idempotency_key: None,
4301 concurrency_key: None,
4302 priority: 0,
4303 concurrency_limits: Vec::new(),
4304 max_cost_usd: None,
4305 worker_tags: Vec::new(),
4306 })
4307 .await
4308 .unwrap()
4309 .into_run();
4310
4311 let result = engine.fail_orphaned_steps(run.id, "timeout").await;
4312 assert!(result.is_ok());
4313 }
4314
4315 #[tokio::test]
4316 async fn fail_orphaned_steps_preserves_existing_error() {
4317 let engine = create_test_engine();
4318 let run = engine
4319 .store()
4320 .create_run(NewRun {
4321 created_by: None,
4322 workflow_name: "test".to_string(),
4323 trigger: TriggerKind::Manual,
4324 payload: json!({}),
4325 max_retries: 0,
4326 handler_version: None,
4327 labels: HashMap::new(),
4328 scheduled_at: None,
4329 idempotency_key: None,
4330 concurrency_key: None,
4331 priority: 0,
4332 concurrency_limits: Vec::new(),
4333 max_cost_usd: None,
4334 worker_tags: Vec::new(),
4335 })
4336 .await
4337 .unwrap()
4338 .into_run();
4339
4340 let step_with_error = create_step_with_status(
4341 engine.store(),
4342 run.id,
4343 "already-errored",
4344 0,
4345 StepStatus::Running,
4346 )
4347 .await;
4348
4349 engine
4350 .store()
4351 .update_step(
4352 step_with_error.id,
4353 StepUpdate {
4354 error: Some("real error from provider".to_string()),
4355 ..StepUpdate::default()
4356 },
4357 )
4358 .await
4359 .unwrap();
4360
4361 let step_no_error = create_step_with_status(
4362 engine.store(),
4363 run.id,
4364 "no-error-yet",
4365 1,
4366 StepStatus::Running,
4367 )
4368 .await;
4369
4370 engine
4371 .fail_orphaned_steps(run.id, "parent run failed")
4372 .await
4373 .unwrap();
4374
4375 let updated_with = engine
4376 .store()
4377 .get_step(step_with_error.id)
4378 .await
4379 .unwrap()
4380 .unwrap();
4381 assert_eq!(updated_with.status.state, StepStatus::Failed);
4382 assert_eq!(
4383 updated_with.error.as_deref(),
4384 Some("real error from provider"),
4385 );
4386
4387 let updated_without = engine
4388 .store()
4389 .get_step(step_no_error.id)
4390 .await
4391 .unwrap()
4392 .unwrap();
4393 assert_eq!(updated_without.status.state, StepStatus::Failed);
4394 assert_eq!(updated_without.error.as_deref(), Some("parent run failed"),);
4395 }
4396}