1use std::collections::{HashMap, HashSet};
10use std::fmt;
11use std::sync::{Arc, Mutex};
12use std::time::Instant;
13
14use chrono::{DateTime, TimeDelta, Utc};
15use rust_decimal::Decimal;
16use serde_json::{Value, to_value};
17use tokio::spawn;
18use tracing::{error, info, warn};
19use uuid::Uuid;
20
21use ironflow_core::error::OperationError;
22#[cfg(feature = "prometheus")]
23use ironflow_core::metric_names::{
24 RUN_BUDGET_EXCEEDED_TOTAL, RUN_COST_USD, RUN_DURATION_SECONDS, RUNS_ACTIVE, RUNS_TOTAL,
25};
26use ironflow_core::provider::{AgentProvider, LABEL_ROOT_RUN_ID};
27use ironflow_store::error::StoreError;
28use ironflow_store::models::{
29 ConcurrencyLimit, LeaseUpdate, NewRun, NewSignal, ProviderKind, Run, RunActor, RunCreation,
30 RunFilter, RunStatus, RunUpdate, SignalInsert, SignalStepResolution, StepStatus, StepUpdate,
31 TriggerKind, normalize_worker_tags, validate_concurrency_limits, validate_worker_tags,
32};
33use ironflow_store::store::Store;
34#[cfg(feature = "prometheus")]
35use metrics::{counter, gauge, histogram};
36
37use crate::artifact::ArtifactSink;
38use crate::budget::{BudgetConfig, month_start};
39use crate::context::{PARENT_RUN_ID_LABEL, WorkflowContext, interrupt_running_steps};
40use crate::error::EngineError;
41use crate::executor::{StepInterceptor, StepResult};
42use crate::guard::{WorkflowGuardConfig, new_shared_guard_state};
43use crate::handler::{WorkflowHandler, WorkflowInfo};
44use crate::log_sender::LogSender;
45use crate::notify::{
46 ApprovalRequestedEvent, Event, EventPublisher, EventSubscriber, RunBudgetExceededEvent,
47 RunFailedEvent, RunStatusChangedEvent, SignalAwaitedEvent, SignalReceivedEvent,
48 WorkflowEventBus,
49};
50use crate::plan::{
51 ExecutionPlan, PlanOptions, PlanRecorder, SharedPlanRecorder, estimate_durations, lock_plan,
52};
53use crate::retry_policy::{backoff_for_retry, is_run_retryable};
54use crate::schedule::CronSchedule;
55use crate::signal::{
56 Signal, SignalDelivery, SignalRejected, SignalResumed, received_output, validate_step_payload,
57};
58use ironflow_core::decision::DecisionProvider;
59
60#[derive(Debug, Clone)]
79pub struct WorkflowResult {
80 pub run: Run,
82 pub steps: Vec<StepResult>,
84}
85
86#[derive(Debug, Clone, Default)]
105pub struct EnqueueOptions {
106 pub max_retries: u32,
108 pub labels: HashMap<String, String>,
110 pub scheduled_at: Option<DateTime<Utc>>,
113 pub max_cost_usd: Option<Decimal>,
117 pub created_by: Option<RunActor>,
120 pub idempotency_key: Option<String>,
126 pub concurrency_key: Option<String>,
132 pub concurrency_limits: Vec<ConcurrencyLimit>,
140 pub worker_tags: Vec<String>,
146}
147
148#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
159pub enum ExecutionMode {
160 #[default]
164 Local,
165 Workers,
168}
169
170pub struct Engine {
211 store: Arc<dyn Store>,
212 provider: Arc<dyn AgentProvider>,
213 handlers: HashMap<String, Arc<dyn WorkflowHandler>>,
214 event_publisher: EventPublisher,
215 log_sender: Option<LogSender>,
216 budget: BudgetConfig,
217 artifact_sink: Option<Arc<dyn ArtifactSink>>,
218 guard_config: Option<WorkflowGuardConfig>,
219 event_bus: Option<WorkflowEventBus>,
220 decision_provider: Option<Arc<dyn DecisionProvider>>,
221 step_interceptor: Option<Arc<dyn StepInterceptor>>,
222 execution_mode: ExecutionMode,
223 worker_tags: Option<Arc<Vec<String>>>,
224}
225
226fn validate_category(handler_name: &str, category: &str) -> Result<(), EngineError> {
236 let reject = |reason: &str| {
237 Err(EngineError::InvalidWorkflow(format!(
238 "handler '{handler_name}' has invalid category '{category}': {reason}"
239 )))
240 };
241
242 if category.is_empty() {
243 return reject("empty category");
244 }
245 if category.starts_with('/') {
246 return reject("leading '/'");
247 }
248 if category.ends_with('/') {
249 return reject("trailing '/'");
250 }
251 for segment in category.split('/') {
252 if segment.is_empty() {
253 return reject("empty segment (double '/')");
254 }
255 if segment.trim().is_empty() {
256 return reject("whitespace-only segment");
257 }
258 }
259 Ok(())
260}
261
262fn chain_label(run: &Run, key: &str) -> Option<Uuid> {
267 if !matches!(run.trigger, TriggerKind::Workflow) {
268 return None;
269 }
270 let id = Uuid::parse_str(run.labels.get(key)?).ok()?;
271 (id != run.id).then_some(id)
272}
273
274pub fn chain_root(run: &Run) -> Option<Uuid> {
294 chain_label(run, LABEL_ROOT_RUN_ID)
295}
296
297fn chain_parent(run: &Run) -> Option<Uuid> {
299 chain_label(run, PARENT_RUN_ID_LABEL)
300}
301
302impl Engine {
303 pub fn new(store: Arc<dyn Store>, provider: Arc<dyn AgentProvider>) -> Self {
319 Self {
320 store,
321 provider,
322 handlers: HashMap::new(),
323 event_publisher: EventPublisher::new(),
324 log_sender: None,
325 budget: BudgetConfig::new(),
326 artifact_sink: None,
327 guard_config: None,
328 event_bus: None,
329 decision_provider: None,
330 step_interceptor: None,
331 execution_mode: ExecutionMode::default(),
332 worker_tags: None,
333 }
334 }
335
336 pub fn with_decision_provider(mut self, provider: Arc<dyn DecisionProvider>) -> Self {
357 self.decision_provider = Some(provider);
358 self
359 }
360
361 pub fn with_step_interceptor(mut self, interceptor: Arc<dyn StepInterceptor>) -> Self {
385 self.step_interceptor = Some(interceptor);
386 self
387 }
388
389 pub fn step_interceptor(&self) -> Option<&Arc<dyn StepInterceptor>> {
391 self.step_interceptor.as_ref()
392 }
393
394 pub fn with_budget_config(mut self, budget: BudgetConfig) -> Self {
415 self.budget = budget;
416 self
417 }
418
419 pub fn budget_config(&self) -> &BudgetConfig {
421 &self.budget
422 }
423
424 pub fn with_guard_config(mut self, config: WorkflowGuardConfig) -> Self {
446 self.guard_config = Some(config);
447 self
448 }
449
450 pub fn guard_config(&self) -> Option<&WorkflowGuardConfig> {
452 self.guard_config.as_ref()
453 }
454
455 pub fn with_execution_mode(mut self, mode: ExecutionMode) -> Self {
476 self.execution_mode = mode;
477 self
478 }
479
480 pub fn execution_mode(&self) -> ExecutionMode {
482 self.execution_mode
483 }
484
485 pub fn set_worker_tags(&mut self, tags: Vec<String>) {
508 self.worker_tags = Some(Arc::new(tags));
509 }
510
511 pub fn worker_tags(&self) -> Option<&[String]> {
529 self.worker_tags.as_deref().map(Vec::as_slice)
530 }
531
532 pub fn set_log_sender(&mut self, sender: LogSender) {
538 self.log_sender = Some(sender);
539 }
540
541 pub fn set_artifact_sink(&mut self, sink: Arc<dyn ArtifactSink>) {
559 self.artifact_sink = Some(sink);
560 }
561
562 pub fn artifact_sink(&self) -> Option<&Arc<dyn ArtifactSink>> {
564 self.artifact_sink.as_ref()
565 }
566
567 pub fn set_event_bus(&mut self, bus: WorkflowEventBus) {
584 self.event_bus = Some(bus);
585 }
586
587 pub fn event_bus(&self) -> Option<&WorkflowEventBus> {
589 self.event_bus.as_ref()
590 }
591
592 pub fn store(&self) -> &Arc<dyn Store> {
594 &self.store
595 }
596
597 pub fn provider(&self) -> &Arc<dyn AgentProvider> {
599 &self.provider
600 }
601
602 fn build_context(&self, run: &Run) -> WorkflowContext {
611 let handlers = self.handlers.clone();
612 let resolver: crate::context::HandlerResolver =
613 Arc::new(move |name: &str| handlers.get(name).cloned());
614 let mut ctx = WorkflowContext::with_handler_resolver(
615 run.id,
616 run.workflow_name.clone(),
617 self.store.clone(),
618 self.provider.clone(),
619 resolver,
620 );
621 ctx.carry_over_run_totals(run.retry_count + 1, run.cost_usd, run.duration_ms);
622 ctx.set_max_cost_usd(run.max_cost_usd);
623 ctx.set_run_created_at(run.created_at);
624 if let Some(ref sender) = self.log_sender {
625 ctx.set_log_sender(sender.clone());
626 }
627 if let Some(ref sink) = self.artifact_sink {
628 ctx.set_artifact_sink(sink.clone());
629 }
630 if let Some(ref bus) = self.event_bus {
631 ctx.set_event_bus(bus.clone());
632 }
633 if let Some(ref provider) = self.decision_provider {
634 ctx.set_decision_provider(provider.clone());
635 }
636 if let Some(ref interceptor) = self.step_interceptor {
637 ctx.set_step_interceptor(interceptor.clone());
638 }
639 if let Some(ref tags) = self.worker_tags {
640 ctx.set_worker_tags(tags.clone());
641 }
642 ctx
643 }
644
645 fn build_context_with_guard(
651 &self,
652 run: &Run,
653 handler: &dyn WorkflowHandler,
654 ) -> WorkflowContext {
655 let mut ctx = self.build_context(run);
656 let guard_config = handler.guard_config().or_else(|| self.guard_config.clone());
657 if let Some(config) = guard_config {
658 ctx.set_guard(config, new_shared_guard_state());
659 }
660 ctx
661 }
662
663 async fn check_monthly_quota(&self, workflow_name: &str) -> Result<(), EngineError> {
674 let Some(limit) = self.budget.monthly_cost_limit_usd else {
675 return Ok(());
676 };
677
678 let stats = self
679 .store
680 .get_stats(RunFilter {
681 created_after: Some(month_start(Utc::now())),
682 ..RunFilter::default()
683 })
684 .await?;
685
686 if stats.total_cost_usd < limit {
687 return Ok(());
688 }
689
690 warn!(
691 workflow = %workflow_name,
692 limit_usd = %limit,
693 spent_usd = %stats.total_cost_usd,
694 "monthly cost quota exhausted, refusing new run"
695 );
696
697 #[cfg(feature = "prometheus")]
698 counter!(
699 RUN_BUDGET_EXCEEDED_TOTAL,
700 "workflow" => workflow_name.to_string(),
701 "scope" => "monthly",
702 )
703 .increment(1);
704
705 Err(EngineError::MonthlyBudgetExceeded {
706 limit_usd: limit,
707 spent_usd: stats.total_cost_usd,
708 })
709 }
710
711 pub fn register(&mut self, handler: impl WorkflowHandler + 'static) -> Result<(), EngineError> {
755 let name = handler.name().to_string();
756 if self.handlers.contains_key(&name) {
757 return Err(EngineError::InvalidWorkflow(format!(
758 "handler '{}' already registered",
759 name
760 )));
761 }
762 if let Some(category) = handler.category() {
763 validate_category(&name, category)?;
764 }
765 self.handlers.insert(name, Arc::new(handler));
766 Ok(())
767 }
768
769 pub fn register_boxed(&mut self, handler: Box<dyn WorkflowHandler>) -> Result<(), EngineError> {
776 let name = handler.name().to_string();
777 if self.handlers.contains_key(&name) {
778 return Err(EngineError::InvalidWorkflow(format!(
779 "handler '{}' already registered",
780 name
781 )));
782 }
783 if let Some(category) = handler.category() {
784 validate_category(&name, category)?;
785 }
786 self.handlers.insert(name, Arc::from(handler));
787 Ok(())
788 }
789
790 pub fn get_handler(&self, name: &str) -> Option<&Arc<dyn WorkflowHandler>> {
792 self.handlers.get(name)
793 }
794
795 pub fn handler_names(&self) -> Vec<&str> {
797 self.handlers.keys().map(|s| s.as_str()).collect()
798 }
799
800 pub fn handler_info(&self, name: &str) -> Option<WorkflowInfo> {
802 self.handlers.get(name).map(|h| h.describe())
803 }
804
805 pub fn scheduled_handlers(&self) -> Vec<(&str, &CronSchedule)> {
829 self.handlers
830 .iter()
831 .filter_map(|(name, handler)| handler.schedule().map(|sched| (name.as_str(), sched)))
832 .collect()
833 }
834
835 pub fn subscribe(
860 &mut self,
861 subscriber: impl EventSubscriber + 'static,
862 event_types: &[&'static str],
863 ) {
864 self.event_publisher.subscribe(subscriber, event_types);
865 }
866
867 pub fn event_publisher(&self) -> &EventPublisher {
872 &self.event_publisher
873 }
874
875 #[tracing::instrument(name = "engine.run_handler", skip_all, fields(workflow = %handler_name))]
905 pub async fn run_handler(
906 &self,
907 handler_name: &str,
908 trigger: TriggerKind,
909 payload: Value,
910 ) -> Result<WorkflowResult, EngineError> {
911 let handler = self
912 .handlers
913 .get(handler_name)
914 .ok_or_else(|| {
915 EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
916 })?
917 .clone();
918
919 self.check_monthly_quota(handler_name).await?;
920
921 let handler_version = handler.version().map(str::to_string);
922 let max_cost_usd = self
923 .budget
924 .resolve_run_cap(None, handler.default_max_cost_usd());
925 let run = self
926 .store
927 .create_run(NewRun {
928 created_by: None,
929 workflow_name: handler_name.to_string(),
930 trigger,
931 payload,
932 max_retries: 0,
933 handler_version,
934 labels: handler.default_labels(),
935 scheduled_at: None,
936 idempotency_key: None,
937 concurrency_key: None,
938 concurrency_limits: Vec::new(),
939 max_cost_usd,
940 worker_tags: normalize_worker_tags(handler.required_worker_tags()),
941 })
942 .await?
943 .into_run();
944
945 let run_id = run.id;
946 info!(run_id = %run_id, handler_version = run.handler_version.as_deref().unwrap_or(""), "run created");
947
948 self.store
949 .update_run_status(run_id, RunStatus::Running)
950 .await?;
951
952 #[cfg(feature = "prometheus")]
953 gauge!(RUNS_ACTIVE, "workflow" => handler_name.to_string()).increment(1.0);
954
955 let run_start = Instant::now();
956 let mut ctx = self.build_context_with_guard(&run, handler.as_ref());
957
958 let result = handler.execute(&mut ctx).await;
959 self.finalize_run(run_id, handler_name, result, &ctx, run_start, run.labels)
960 .await
961 }
962
963 #[tracing::instrument(name = "engine.plan_handler", skip_all, fields(workflow = %handler_name))]
1005 pub async fn plan_handler(
1006 &self,
1007 handler_name: &str,
1008 payload: Value,
1009 options: PlanOptions,
1010 ) -> Result<ExecutionPlan, EngineError> {
1011 if options.max_depth == 0 {
1012 return Err(EngineError::InvalidWorkflow(
1013 "max_depth must be at least 1".to_string(),
1014 ));
1015 }
1016
1017 let handler = self
1018 .handlers
1019 .get(handler_name)
1020 .ok_or_else(|| {
1021 EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
1022 })?
1023 .clone();
1024
1025 let estimates = if options.estimate_durations {
1026 estimate_durations(&self.store, handler_name, options.sample_runs).await?
1027 } else {
1028 HashMap::new()
1029 };
1030
1031 let shared: SharedPlanRecorder = Arc::new(Mutex::new(PlanRecorder::new(
1032 handler_name.to_string(),
1033 payload,
1034 options.max_depth,
1035 estimates,
1036 )));
1037
1038 let handlers = self.handlers.clone();
1041 let resolver: crate::context::HandlerResolver =
1042 Arc::new(move |name: &str| handlers.get(name).cloned());
1043 let mut ctx = WorkflowContext::with_handler_resolver(
1044 Uuid::now_v7(),
1045 handler_name.to_string(),
1046 self.store.clone(),
1047 self.provider.clone(),
1048 resolver,
1049 );
1050 ctx.set_plan(shared.clone());
1051
1052 if let Err(err) = handler.execute(&mut ctx).await {
1053 lock_plan(&shared).fail(err.to_string());
1054 }
1055 drop(ctx);
1056
1057 let plan = match Arc::try_unwrap(shared) {
1058 Ok(mutex) => mutex
1059 .into_inner()
1060 .unwrap_or_else(|poisoned| poisoned.into_inner())
1061 .into_plan(),
1062 Err(shared) => lock_plan(&shared).snapshot(),
1063 };
1064
1065 info!(
1066 workflow = %handler_name,
1067 steps = plan.steps.len(),
1068 truncated = plan.truncated,
1069 "execution plan built"
1070 );
1071
1072 Ok(plan)
1073 }
1074
1075 #[tracing::instrument(name = "engine.enqueue_handler", skip_all, fields(workflow = %handler_name))]
1086 pub async fn enqueue_handler(
1087 &self,
1088 handler_name: &str,
1089 trigger: TriggerKind,
1090 payload: Value,
1091 max_retries: u32,
1092 ) -> Result<Run, EngineError> {
1093 self.enqueue_handler_with_options(
1094 handler_name,
1095 trigger,
1096 payload,
1097 EnqueueOptions {
1098 max_retries,
1099 ..Default::default()
1100 },
1101 )
1102 .await
1103 .map(RunCreation::into_run)
1104 }
1105
1106 #[tracing::instrument(name = "engine.enqueue_handler_with_options", skip_all, fields(workflow = %handler_name))]
1156 pub async fn enqueue_handler_with_options(
1157 &self,
1158 handler_name: &str,
1159 trigger: TriggerKind,
1160 payload: Value,
1161 options: EnqueueOptions,
1162 ) -> Result<RunCreation, EngineError> {
1163 let EnqueueOptions {
1164 max_retries,
1165 labels,
1166 scheduled_at,
1167 max_cost_usd,
1168 created_by,
1169 idempotency_key,
1170 concurrency_key,
1171 concurrency_limits,
1172 worker_tags,
1173 } = options;
1174
1175 validate_concurrency_limits(&concurrency_limits)
1178 .map_err(EngineError::InvalidConcurrencyLimit)?;
1179 validate_worker_tags(&worker_tags).map_err(EngineError::InvalidWorkerTag)?;
1180
1181 let handler = self.handlers.get(handler_name).ok_or_else(|| {
1182 EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
1183 })?;
1184
1185 self.check_monthly_quota(handler_name).await?;
1186
1187 let handler_version = handler.version().map(str::to_string);
1188 let mut merged_labels = handler.default_labels();
1189 merged_labels.extend(labels);
1190 let resolved_cap = self
1191 .budget
1192 .resolve_run_cap(max_cost_usd, handler.default_max_cost_usd());
1193 let required_tags = normalize_worker_tags(
1194 handler
1195 .required_worker_tags()
1196 .into_iter()
1197 .chain(worker_tags),
1198 );
1199
1200 let creation = self
1201 .store
1202 .create_run(NewRun {
1203 workflow_name: handler_name.to_string(),
1204 trigger,
1205 payload,
1206 max_retries,
1207 handler_version,
1208 labels: merged_labels,
1209 scheduled_at,
1210 created_by,
1211 idempotency_key,
1212 concurrency_key,
1213 concurrency_limits,
1214 max_cost_usd: resolved_cap,
1215 worker_tags: required_tags,
1216 })
1217 .await?;
1218
1219 match &creation {
1220 RunCreation::Created(run) => info!(
1221 run_id = %run.id,
1222 workflow = %handler_name,
1223 max_cost_usd = ?resolved_cap,
1224 "handler run enqueued"
1225 ),
1226 RunCreation::Existing(run) => info!(
1227 run_id = %run.id,
1228 workflow = %handler_name,
1229 "idempotent replay, nothing enqueued"
1230 ),
1231 }
1232
1233 Ok(creation)
1234 }
1235
1236 #[tracing::instrument(name = "engine.execute_handler_run", skip_all, fields(run_id = %run_id))]
1256 pub async fn execute_handler_run(&self, run_id: Uuid) -> Result<WorkflowResult, EngineError> {
1257 let run = self
1258 .store
1259 .get_run(run_id)
1260 .await?
1261 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
1262
1263 if let Some(root_run_id) = chain_root(&run) {
1264 return self.resume_chain(run, root_run_id).await;
1265 }
1266
1267 let handler = self
1268 .handlers
1269 .get(&run.workflow_name)
1270 .ok_or_else(|| {
1271 EngineError::InvalidWorkflow(format!(
1272 "no handler registered: {}",
1273 run.workflow_name
1274 ))
1275 })?
1276 .clone();
1277
1278 #[cfg(feature = "prometheus")]
1279 gauge!(RUNS_ACTIVE, "workflow" => run.workflow_name.clone()).increment(1.0);
1280
1281 let run_start = Instant::now();
1282 let mut ctx = self.build_context_with_guard(&run, handler.as_ref());
1283
1284 let result = if handler.is_version_compatible(run.handler_version.as_deref()) {
1297 ctx.load_replay_steps().await?;
1298 self.release_then_execute(run_id, handler.as_ref(), &mut ctx)
1299 .await
1300 } else {
1301 Err(EngineError::HandlerVersionMismatch {
1302 run_id,
1303 workflow_name: run.workflow_name.clone(),
1304 run_version: run
1305 .handler_version
1306 .clone()
1307 .unwrap_or_else(|| "unknown".to_string()),
1308 current_version: handler
1309 .version()
1310 .map(str::to_string)
1311 .unwrap_or_else(|| "unknown".to_string()),
1312 })
1313 };
1314
1315 self.finalize_run(
1316 run_id,
1317 &run.workflow_name,
1318 result,
1319 &ctx,
1320 run_start,
1321 run.labels,
1322 )
1323 .await
1324 }
1325
1326 #[tracing::instrument(name = "engine.execute_run", skip_all, fields(run_id = %run_id))]
1334 pub async fn execute_run(&self, run_id: Uuid) -> Result<WorkflowResult, EngineError> {
1335 self.execute_handler_run(run_id).await
1336 }
1337
1338 #[tracing::instrument(name = "engine.resume_run", skip_all, fields(run_id = %run_id))]
1367 pub async fn resume_run(&self, run_id: Uuid) -> Result<WorkflowResult, EngineError> {
1368 let run = self
1369 .store
1370 .get_run(run_id)
1371 .await?
1372 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
1373
1374 if let Some(root_run_id) = chain_root(&run) {
1375 return self.resume_chain(run, root_run_id).await;
1376 }
1377
1378 self.resume_loaded_run(run).await
1379 }
1380
1381 async fn resume_chain(
1396 &self,
1397 child: Run,
1398 root_run_id: Uuid,
1399 ) -> Result<WorkflowResult, EngineError> {
1400 let child_run_id = child.id;
1401 let lease = child.worker_id.zip(child.lease_expires_at);
1402 let root = self
1403 .store
1404 .get_run(root_run_id)
1405 .await?
1406 .ok_or(EngineError::Store(StoreError::RunNotFound(root_run_id)))?;
1407
1408 match root.status.state {
1409 RunStatus::AwaitingApproval | RunStatus::Pending => {
1410 self.move_root_to_running(root_run_id, lease.as_ref())
1411 .await?;
1412 }
1413 RunStatus::Sleeping => {
1414 self.store
1415 .update_run_status(root_run_id, RunStatus::Pending)
1416 .await?;
1417 self.move_root_to_running(root_run_id, lease.as_ref())
1418 .await?;
1419 }
1420 other => {
1421 let reason = format!(
1422 "cannot resume child run {child_run_id}: root run {root_run_id} is {other}"
1423 );
1424 if let Err(err) = self
1425 .fail_or_schedule_retry(child_run_id, &reason, false, None, None)
1426 .await
1427 {
1428 error!(
1429 run_id = %child_run_id,
1430 error = %err,
1431 "failed to fail a child run whose root cannot resume"
1432 );
1433 }
1434 return Err(EngineError::InvalidWorkflow(reason));
1435 }
1436 }
1437
1438 if lease.is_some() {
1439 self.store
1440 .update_run(
1441 child_run_id,
1442 RunUpdate {
1443 lease: Some(LeaseUpdate::Release),
1444 ..RunUpdate::default()
1445 },
1446 )
1447 .await?;
1448 }
1449
1450 info!(
1451 run_id = %child_run_id,
1452 root_run_id = %root_run_id,
1453 lease_transferred = lease.is_some(),
1454 "child run resumed through its root run"
1455 );
1456
1457 let root = self
1458 .store
1459 .get_run(root_run_id)
1460 .await?
1461 .ok_or(EngineError::Store(StoreError::RunNotFound(root_run_id)))?;
1462 self.resume_loaded_run(root).await
1463 }
1464
1465 async fn move_root_to_running(
1471 &self,
1472 root_run_id: Uuid,
1473 lease: Option<&(String, DateTime<Utc>)>,
1474 ) -> Result<(), EngineError> {
1475 match lease {
1476 Some((worker_id, expires_at)) => {
1477 self.store
1478 .update_run(
1479 root_run_id,
1480 RunUpdate {
1481 status: Some(RunStatus::Running),
1482 lease: Some(LeaseUpdate::Set {
1483 worker_id: worker_id.clone(),
1484 expires_at: *expires_at,
1485 }),
1486 ..RunUpdate::default()
1487 },
1488 )
1489 .await?;
1490 }
1491 None => {
1492 self.store
1493 .update_run_status(root_run_id, RunStatus::Running)
1494 .await?;
1495 }
1496 }
1497 Ok(())
1498 }
1499
1500 async fn resume_loaded_run(&self, run: Run) -> Result<WorkflowResult, EngineError> {
1502 let run_id = run.id;
1503 let handler = self
1504 .handlers
1505 .get(&run.workflow_name)
1506 .ok_or_else(|| {
1507 EngineError::InvalidWorkflow(format!(
1508 "no handler registered: {}",
1509 run.workflow_name
1510 ))
1511 })?
1512 .clone();
1513
1514 info!(run_id = %run_id, workflow = %run.workflow_name, "resuming run after approval");
1515
1516 let run_start = Instant::now();
1517 let mut ctx = self.build_context_with_guard(&run, handler.as_ref());
1518
1519 let result = if handler.is_version_compatible(run.handler_version.as_deref()) {
1520 ctx.load_replay_steps().await?;
1521 self.release_then_execute(run_id, handler.as_ref(), &mut ctx)
1522 .await
1523 } else {
1524 Err(EngineError::HandlerVersionMismatch {
1525 run_id,
1526 workflow_name: run.workflow_name.clone(),
1527 run_version: run
1528 .handler_version
1529 .clone()
1530 .unwrap_or_else(|| "unknown".to_string()),
1531 current_version: handler
1532 .version()
1533 .map(str::to_string)
1534 .unwrap_or_else(|| "unknown".to_string()),
1535 })
1536 };
1537
1538 self.finalize_run(
1539 run_id,
1540 &run.workflow_name,
1541 result,
1542 &ctx,
1543 run_start,
1544 run.labels,
1545 )
1546 .await
1547 }
1548
1549 pub async fn deliver_signal(
1595 self: &Arc<Self>,
1596 signal: NewSignal,
1597 ) -> Result<SignalDelivery, EngineError> {
1598 if signal.name.trim().is_empty() {
1599 return Err(EngineError::InvalidSignal(
1600 "signal name must not be empty".to_string(),
1601 ));
1602 }
1603 if signal.key.trim().is_empty() {
1604 return Err(EngineError::InvalidSignal(
1605 "signal key must not be empty".to_string(),
1606 ));
1607 }
1608
1609 let stored = match self.store.insert_signal(signal).await? {
1610 SignalInsert::Created(stored) => stored,
1611 SignalInsert::Duplicate(existing) => {
1612 info!(
1613 signal_id = %existing.id,
1614 signal = %existing.name,
1615 key = %existing.key,
1616 "duplicate signal ignored"
1617 );
1618 return Ok(SignalDelivery {
1619 signal_id: existing.id,
1620 duplicate: true,
1621 resumed: Vec::new(),
1622 rejected: Vec::new(),
1623 });
1624 }
1625 };
1626
1627 let waiters = self
1628 .store
1629 .list_signal_waiters(&stored.name, &stored.key)
1630 .await?;
1631 let mut resumed = Vec::new();
1632 let mut rejected = Vec::new();
1633
1634 for step in waiters {
1635 if let Err(error) = validate_step_payload(step.input.as_ref(), &stored.payload) {
1636 rejected.push(SignalRejected {
1637 run_id: step.run_id,
1638 step_id: step.id,
1639 error,
1640 });
1641 continue;
1642 }
1643
1644 match self
1645 .store
1646 .resolve_signal_step(step.id, received_output(&stored))
1647 .await
1648 {
1649 Ok(SignalStepResolution::Resolved {
1650 run_id,
1651 run_resumed,
1652 }) => {
1653 resumed.push(SignalResumed {
1654 run_id,
1655 step_id: step.id,
1656 });
1657 if run_resumed && self.execution_mode == ExecutionMode::Local {
1658 self.spawn_local_resume(run_id);
1659 }
1660 }
1661 Ok(SignalStepResolution::NotWaiting { .. }) => {}
1663 Err(err) => {
1664 error!(
1665 run_id = %step.run_id,
1666 step_id = %step.id,
1667 error = %err,
1668 "failed to resolve a waiting signal step"
1669 );
1670 rejected.push(SignalRejected {
1671 run_id: step.run_id,
1672 step_id: step.id,
1673 error: err.to_string(),
1674 });
1675 }
1676 }
1677 }
1678
1679 info!(
1680 signal_id = %stored.id,
1681 signal = %stored.name,
1682 key = %stored.key,
1683 resumed = resumed.len(),
1684 rejected = rejected.len(),
1685 "signal received"
1686 );
1687 self.event_publisher
1688 .publish(Event::SignalReceived(SignalReceivedEvent {
1689 signal_id: stored.id,
1690 name: stored.name.clone(),
1691 key: stored.key.clone(),
1692 resumed_runs: resumed.iter().map(|r| r.run_id).collect(),
1693 at: stored.received_at,
1694 }));
1695
1696 Ok(SignalDelivery {
1697 signal_id: stored.id,
1698 duplicate: false,
1699 resumed,
1700 rejected,
1701 })
1702 }
1703
1704 pub async fn send_signal<S: Signal>(
1742 self: &Arc<Self>,
1743 signal: &S,
1744 key: &str,
1745 idempotency_id: Option<&str>,
1746 ) -> Result<SignalDelivery, EngineError> {
1747 let payload = to_value(signal)?;
1748 self.deliver_signal(NewSignal {
1749 name: S::NAME.to_string(),
1750 key: key.to_string(),
1751 payload,
1752 idempotency_id: idempotency_id.map(str::to_string),
1753 })
1754 .await
1755 }
1756
1757 pub(crate) fn spawn_local_resume(self: &Arc<Self>, run_id: Uuid) {
1763 let engine = Arc::clone(self);
1764 spawn(async move {
1765 if let Err(err) = engine
1766 .store
1767 .update_run_status(run_id, RunStatus::Running)
1768 .await
1769 {
1770 error!(run_id = %run_id, error = %err, "failed to restart a woken run");
1771 return;
1772 }
1773 if let Err(err) = engine.resume_run(run_id).await {
1774 error!(run_id = %run_id, error = %err, "failed to resume a woken run");
1775 }
1776 });
1777 }
1778
1779 pub async fn fail_or_schedule_retry(
1827 &self,
1828 run_id: Uuid,
1829 error: &str,
1830 retryable: bool,
1831 cost_usd: Option<Decimal>,
1832 duration_ms: Option<u64>,
1833 ) -> Result<RunStatus, EngineError> {
1834 let run = self
1835 .store
1836 .get_run(run_id)
1837 .await?
1838 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
1839
1840 let has_attempts_left = run.retry_count < run.max_retries;
1841 let update = if retryable && has_attempts_left {
1842 let backoff = backoff_for_retry(run.retry_count);
1843 let scheduled_at = Utc::now() + TimeDelta::milliseconds(backoff.as_millis() as i64);
1844
1845 info!(
1846 run_id = %run_id,
1847 workflow = %run.workflow_name,
1848 attempt = run.retry_count + 1,
1849 max_retries = run.max_retries,
1850 backoff_secs = backoff.as_secs(),
1851 scheduled_at = %scheduled_at,
1852 "run failed, scheduling retry"
1853 );
1854
1855 RunUpdate {
1856 status: Some(RunStatus::Retrying),
1857 error: Some(error.to_string()),
1858 increment_retry: true,
1859 cost_usd,
1860 duration_ms,
1861 scheduled_at: Some(scheduled_at),
1862 ..RunUpdate::default()
1863 }
1864 } else {
1865 RunUpdate {
1866 status: Some(RunStatus::Failed),
1867 error: Some(error.to_string()),
1868 cost_usd,
1869 duration_ms,
1870 completed_at: Some(Utc::now()),
1871 ..RunUpdate::default()
1872 }
1873 };
1874
1875 let status = update.status.unwrap_or(RunStatus::Failed);
1876 self.store.update_run(run_id, update).await?;
1877 self.fail_orphaned_steps(run_id, error).await?;
1878 self.cancel_descendants_of_stopped_run(run_id, error).await;
1881
1882 Ok(status)
1883 }
1884
1885 pub async fn interrupt_running_steps(&self, run_id: Uuid) -> Result<(), EngineError> {
1915 interrupt_running_steps(self.store.as_ref(), run_id).await
1916 }
1917
1918 pub async fn fail_orphaned_steps(
1932 &self,
1933 run_id: Uuid,
1934 error_message: &str,
1935 ) -> Result<(), EngineError> {
1936 let steps = self.store.list_steps(run_id).await?;
1937 let now = Utc::now();
1938
1939 for step in steps {
1940 if step.status.state.is_terminal() {
1941 continue;
1942 }
1943
1944 let (target_status, error) = match step.status.state {
1945 StepStatus::Running | StepStatus::AwaitingApproval => {
1946 let err = if step.error.is_some() {
1947 None
1948 } else {
1949 Some(error_message.to_string())
1950 };
1951 (StepStatus::Failed, err)
1952 }
1953 StepStatus::Pending => (StepStatus::Skipped, None),
1954 _ => continue,
1955 };
1956
1957 if let Err(e) = self
1958 .store
1959 .update_step(
1960 step.id,
1961 StepUpdate {
1962 status: Some(target_status),
1963 error,
1964 completed_at: Some(now),
1965 ..StepUpdate::default()
1966 },
1967 )
1968 .await
1969 {
1970 warn!(
1971 run_id = %run_id,
1972 step_id = %step.id,
1973 step_name = %step.name,
1974 error = %e,
1975 "failed to cleanup orphaned step"
1976 );
1977 } else {
1978 info!(
1979 run_id = %run_id,
1980 step_id = %step.id,
1981 step_name = %step.name,
1982 from = %step.status.state,
1983 to = %target_status,
1984 "cleaned up orphaned step"
1985 );
1986 }
1987 }
1988
1989 Ok(())
1990 }
1991
1992 async fn release_then_execute(
1997 &self,
1998 run_id: Uuid,
1999 handler: &dyn WorkflowHandler,
2000 ctx: &mut WorkflowContext,
2001 ) -> Result<(), EngineError> {
2002 match self.provider.release_run(&run_id.to_string()).await {
2003 Ok(()) => handler.execute(ctx).await,
2004 Err(e) => Err(EngineError::Operation(OperationError::Agent(e))),
2005 }
2006 }
2007
2008 async fn finalize_run(
2014 &self,
2015 run_id: Uuid,
2016 workflow_name: &str,
2017 result: Result<(), EngineError>,
2018 ctx: &WorkflowContext,
2019 run_start: Instant,
2020 run_labels: HashMap<String, String>,
2021 ) -> Result<WorkflowResult, EngineError> {
2022 let total_duration = ctx.carried_duration_ms() + run_start.elapsed().as_millis() as u64;
2025 let completed_at = Utc::now();
2026
2027 let final_status;
2028 let final_run;
2029
2030 match result {
2031 Ok(()) => {
2032 final_status = if ctx.has_allowed_failure() {
2033 RunStatus::Warning
2034 } else {
2035 RunStatus::Completed
2036 };
2037 final_run = self
2038 .store
2039 .update_run_returning(
2040 run_id,
2041 RunUpdate {
2042 status: Some(final_status),
2043 cost_usd: Some(ctx.total_cost_usd()),
2044 duration_ms: Some(total_duration),
2045 completed_at: Some(completed_at),
2046 output: ctx.output().cloned(),
2047 ..RunUpdate::default()
2048 },
2049 )
2050 .await?;
2051
2052 info!(
2053 run_id = %run_id,
2054 status = %final_status,
2055 cost_usd = %ctx.total_cost_usd(),
2056 duration_ms = total_duration,
2057 "run completed"
2058 );
2059 }
2060 Err(EngineError::ApprovalRequired {
2061 run_id: approval_run_id,
2062 step_id,
2063 ref message,
2064 }) => {
2065 final_status = RunStatus::AwaitingApproval;
2066 final_run = self
2067 .store
2068 .update_run_returning(
2069 run_id,
2070 RunUpdate {
2071 status: Some(RunStatus::AwaitingApproval),
2072 cost_usd: Some(ctx.total_cost_usd()),
2073 duration_ms: Some(total_duration),
2074 ..RunUpdate::default()
2075 },
2076 )
2077 .await?;
2078
2079 info!(
2080 run_id = %approval_run_id,
2081 step_id = %step_id,
2082 message = %message,
2083 "run awaiting approval"
2084 );
2085
2086 self.publish_approval_requested(approval_run_id, step_id, message)
2087 .await?;
2088 }
2089 Err(EngineError::ChildSuspended {
2090 run_id: child_run_id,
2091 ref cause,
2092 }) => {
2093 final_status = cause.suspension_status();
2094 final_run = self
2098 .store
2099 .update_run_returning(
2100 run_id,
2101 RunUpdate {
2102 status: Some(final_status),
2103 cost_usd: Some(ctx.total_cost_usd()),
2104 duration_ms: Some(total_duration),
2105 ..RunUpdate::default()
2106 },
2107 )
2108 .await?;
2109
2110 let leaf = cause.suspension_leaf();
2111 info!(
2112 run_id = %run_id,
2113 child_run_id = %child_run_id,
2114 status = %final_status,
2115 cause = %leaf,
2116 "run suspended with its child run"
2117 );
2118
2119 match leaf {
2120 EngineError::ApprovalRequired {
2121 run_id: approval_run_id,
2122 step_id,
2123 message,
2124 } => {
2125 self.publish_approval_requested(*approval_run_id, *step_id, message)
2126 .await?;
2127 }
2128 EngineError::SignalWaiting {
2129 run_id: wait_run_id,
2130 step_id,
2131 step_name,
2132 name,
2133 key,
2134 deadline_at,
2135 } => {
2136 self.event_publisher
2137 .publish(Event::SignalAwaited(SignalAwaitedEvent {
2138 run_id: *wait_run_id,
2139 step_id: *step_id,
2140 step_name: step_name.clone(),
2141 name: name.clone(),
2142 key: key.clone(),
2143 deadline_at: *deadline_at,
2144 at: Utc::now(),
2145 }));
2146 }
2147 _ => {}
2150 }
2151 }
2152 Err(EngineError::HumanInputRequired {
2153 run_id: input_run_id,
2154 step_id,
2155 ref message,
2156 }) => {
2157 final_status = RunStatus::AwaitingApproval;
2158 final_run = self
2159 .store
2160 .update_run_returning(
2161 run_id,
2162 RunUpdate {
2163 status: Some(RunStatus::AwaitingApproval),
2164 cost_usd: Some(ctx.total_cost_usd()),
2165 duration_ms: Some(total_duration),
2166 ..RunUpdate::default()
2167 },
2168 )
2169 .await?;
2170
2171 info!(
2173 run_id = %input_run_id,
2174 step_id = %step_id,
2175 message = %message,
2176 "run awaiting human input"
2177 );
2178 }
2179 Err(EngineError::DelaySleeping {
2180 run_id: delay_run_id,
2181 step_id,
2182 wake_at,
2183 }) => {
2184 final_status = RunStatus::Sleeping;
2185 final_run = self
2186 .store
2187 .update_run_returning(
2188 run_id,
2189 RunUpdate {
2190 status: Some(RunStatus::Sleeping),
2191 cost_usd: Some(ctx.total_cost_usd()),
2192 duration_ms: Some(total_duration),
2193 scheduled_at: Some(wake_at),
2194 ..RunUpdate::default()
2195 },
2196 )
2197 .await?;
2198
2199 info!(
2200 run_id = %delay_run_id,
2201 step_id = %step_id,
2202 wake_at = %wake_at,
2203 "run sleeping until delay elapses"
2204 );
2205 }
2206 Err(EngineError::CapacitySleeping {
2207 run_id: capacity_run_id,
2208 step_id,
2209 ref kind,
2210 wake_at,
2211 }) => {
2212 final_status = RunStatus::Sleeping;
2213 final_run = self
2214 .store
2215 .update_run_returning(
2216 run_id,
2217 RunUpdate {
2218 status: Some(RunStatus::Sleeping),
2219 cost_usd: Some(ctx.total_cost_usd()),
2220 duration_ms: Some(total_duration),
2221 scheduled_at: Some(wake_at),
2222 capacity_wait_kind: Some(ProviderKind::new(kind.as_str())),
2223 ..RunUpdate::default()
2224 },
2225 )
2226 .await?;
2227
2228 info!(
2229 run_id = %capacity_run_id,
2230 step_id = %step_id,
2231 kind = %kind,
2232 wake_at = %wake_at,
2233 "run sleeping until provider capacity returns"
2234 );
2235 }
2236 Err(EngineError::SignalWaiting {
2237 run_id: wait_run_id,
2238 step_id,
2239 ref step_name,
2240 ref name,
2241 ref key,
2242 deadline_at,
2243 }) => {
2244 final_status = RunStatus::Sleeping;
2245 let waiting = self
2249 .store
2250 .suspend_run_on_signal(run_id, step_id, deadline_at)
2251 .await?;
2252 final_run = self
2253 .store
2254 .update_run_returning(
2255 run_id,
2256 RunUpdate {
2257 cost_usd: Some(ctx.total_cost_usd()),
2258 duration_ms: Some(total_duration),
2259 ..RunUpdate::default()
2260 },
2261 )
2262 .await?;
2263
2264 if waiting {
2265 self.event_publisher
2266 .publish(Event::SignalAwaited(SignalAwaitedEvent {
2267 run_id: wait_run_id,
2268 step_id,
2269 step_name: step_name.clone(),
2270 name: name.clone(),
2271 key: key.clone(),
2272 deadline_at,
2273 at: Utc::now(),
2274 }));
2275 }
2276
2277 info!(
2278 run_id = %wait_run_id,
2279 step_id = %step_id,
2280 signal = %name,
2281 key = %key,
2282 deadline_at = %deadline_at,
2283 waiting,
2284 "run sleeping until a signal arrives"
2285 );
2286 }
2287 Err(err) => {
2288 let guardrail_stop = matches!(
2292 err,
2293 EngineError::RunBudgetExceeded { .. } | EngineError::WorkflowGuardRejected(_)
2294 );
2295
2296 final_status = if guardrail_stop {
2297 if let Err(store_err) = self
2298 .store
2299 .update_run(
2300 run_id,
2301 RunUpdate {
2302 status: Some(RunStatus::Cancelled),
2303 error: Some(err.to_string()),
2304 cost_usd: Some(ctx.total_cost_usd()),
2305 duration_ms: Some(total_duration),
2306 completed_at: Some(completed_at),
2307 output: ctx.output().cloned(),
2308 ..RunUpdate::default()
2309 },
2310 )
2311 .await
2312 {
2313 error!(run_id = %run_id, store_error = %store_err, "failed to persist run cancellation");
2314 }
2315 if let Err(cleanup_err) = self
2316 .fail_orphaned_steps(run_id, "run stopped: guardrail limit reached")
2317 .await
2318 {
2319 error!(run_id = %run_id, store_error = %cleanup_err, "failed to cleanup orphaned steps");
2320 }
2321 RunStatus::Cancelled
2322 } else {
2323 if let Some(output) = ctx.output()
2326 && let Err(store_err) = self
2327 .store
2328 .update_run(
2329 run_id,
2330 RunUpdate {
2331 output: Some(output.clone()),
2332 ..RunUpdate::default()
2333 },
2334 )
2335 .await
2336 {
2337 error!(run_id = %run_id, store_error = %store_err, "failed to persist run output");
2338 }
2339 self.fail_or_schedule_retry(
2340 run_id,
2341 &err.to_string(),
2342 is_run_retryable(&err),
2343 Some(ctx.total_cost_usd()),
2344 Some(total_duration),
2345 )
2346 .await
2347 .unwrap_or_else(|store_err| {
2348 error!(run_id = %run_id, store_error = %store_err, "failed to persist run failure");
2349 RunStatus::Failed
2350 })
2351 };
2352
2353 if matches!(err, EngineError::RunBudgetExceeded { .. }) {
2354 self.on_run_budget_exceeded(workflow_name, run_id, &err);
2355 }
2356
2357 error!(run_id = %run_id, status = %final_status, error = %err, "run stopped");
2358
2359 self.publish_run_status_changed(
2360 workflow_name,
2361 run_id,
2362 final_status,
2363 Some(err.to_string()),
2364 ctx,
2365 total_duration,
2366 run_labels,
2367 );
2368
2369 #[cfg(feature = "prometheus")]
2370 self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
2371
2372 return Err(err);
2373 }
2374 }
2375
2376 self.publish_run_status_changed(
2377 workflow_name,
2378 run_id,
2379 final_status,
2380 None,
2381 ctx,
2382 total_duration,
2383 run_labels,
2384 );
2385
2386 #[cfg(feature = "prometheus")]
2387 self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
2388
2389 Ok(WorkflowResult {
2390 run: final_run,
2391 steps: ctx.step_results().to_vec(),
2392 })
2393 }
2394
2395 async fn publish_approval_requested(
2398 &self,
2399 run_id: Uuid,
2400 step_id: Uuid,
2401 message: &str,
2402 ) -> Result<(), EngineError> {
2403 let requirement = self
2404 .store
2405 .get_step(step_id)
2406 .await?
2407 .and_then(|s| s.approval_requirement);
2408 self.event_publisher
2409 .publish(Event::ApprovalRequested(ApprovalRequestedEvent {
2410 run_id,
2411 step_id,
2412 message: message.to_string(),
2413 requirement,
2414 at: Utc::now(),
2415 }));
2416 Ok(())
2417 }
2418
2419 pub async fn fail_ancestors(&self, run_id: Uuid, reason: &str) -> Result<(), EngineError> {
2448 let mut current = self
2449 .store
2450 .get_run(run_id)
2451 .await?
2452 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
2453 let mut visited = HashSet::from([run_id]);
2455
2456 while let Some(parent_id) = chain_parent(¤t) {
2457 if !visited.insert(parent_id) {
2458 break;
2459 }
2460 let status = self
2461 .fail_or_schedule_retry(parent_id, reason, false, None, None)
2462 .await?;
2463 info!(
2464 run_id = %run_id,
2465 ancestor_run_id = %parent_id,
2466 status = %status,
2467 "ancestor run failed with its child"
2468 );
2469 current = self
2470 .store
2471 .get_run(parent_id)
2472 .await?
2473 .ok_or(EngineError::Store(StoreError::RunNotFound(parent_id)))?;
2474 }
2475
2476 Ok(())
2477 }
2478
2479 #[cfg(feature = "prometheus")]
2481 fn emit_run_metrics(
2482 &self,
2483 workflow_name: &str,
2484 status: RunStatus,
2485 duration_ms: u64,
2486 ctx: &WorkflowContext,
2487 ) {
2488 let status_str = status.to_string();
2489 let wf = workflow_name.to_string();
2490
2491 counter!(RUNS_TOTAL, "workflow" => wf.clone(), "status" => status_str.clone()).increment(1);
2492 histogram!(RUN_DURATION_SECONDS, "workflow" => wf.clone(), "status" => status_str)
2493 .record(duration_ms as f64 / 1000.0);
2494 histogram!(RUN_COST_USD, "workflow" => wf.clone()).record(
2495 ctx.total_cost_usd()
2496 .to_string()
2497 .parse::<f64>()
2498 .unwrap_or(0.0),
2499 );
2500 gauge!(RUNS_ACTIVE, "workflow" => wf).decrement(1.0);
2501 }
2502
2503 fn on_run_budget_exceeded(&self, workflow_name: &str, run_id: Uuid, err: &EngineError) {
2509 let EngineError::RunBudgetExceeded {
2510 limit_usd,
2511 spent_usd,
2512 step_budget_usd,
2513 ..
2514 } = err
2515 else {
2516 return;
2517 };
2518
2519 #[cfg(feature = "prometheus")]
2520 counter!(
2521 RUN_BUDGET_EXCEEDED_TOTAL,
2522 "workflow" => workflow_name.to_string(),
2523 "scope" => "run",
2524 )
2525 .increment(1);
2526
2527 self.event_publisher
2528 .publish(Event::RunBudgetExceeded(RunBudgetExceededEvent {
2529 run_id,
2530 workflow_name: workflow_name.to_string(),
2531 limit_usd: *limit_usd,
2532 spent_usd: *spent_usd,
2533 step_budget_usd: *step_budget_usd,
2534 at: Utc::now(),
2535 }));
2536 }
2537
2538 #[allow(clippy::too_many_arguments)]
2543 fn publish_run_status_changed(
2544 &self,
2545 workflow_name: &str,
2546 run_id: Uuid,
2547 to: RunStatus,
2548 error: Option<String>,
2549 ctx: &WorkflowContext,
2550 duration_ms: u64,
2551 labels: HashMap<String, String>,
2552 ) {
2553 let now = Utc::now();
2554 let cost_usd = ctx.total_cost_usd();
2555 let wf = workflow_name.to_string();
2556
2557 self.event_publisher
2558 .publish(Event::RunStatusChanged(RunStatusChangedEvent {
2559 run_id,
2560 workflow_name: wf.clone(),
2561 from: RunStatus::Running,
2562 to,
2563 error: error.clone(),
2564 cost_usd,
2565 duration_ms,
2566 labels: labels.clone(),
2567 at: now,
2568 }));
2569
2570 if to == RunStatus::Failed {
2571 self.event_publisher
2572 .publish(Event::RunFailed(RunFailedEvent {
2573 run_id,
2574 workflow_name: wf,
2575 error,
2576 cost_usd,
2577 duration_ms,
2578 labels,
2579 at: now,
2580 }));
2581 }
2582 }
2583}
2584
2585impl fmt::Debug for Engine {
2586 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
2587 f.debug_struct("Engine")
2588 .field("handlers", &self.handlers.keys().collect::<Vec<_>>())
2589 .finish_non_exhaustive()
2590 }
2591}
2592
2593#[cfg(test)]
2594mod tests {
2595 use super::*;
2596 use crate::config::ShellConfig;
2597 use crate::handler::{HandlerFuture, WorkflowHandler};
2598 use ironflow_core::providers::claude::ClaudeCodeProvider;
2599 use ironflow_core::providers::record_replay::RecordReplayProvider;
2600 use ironflow_store::memory::InMemoryStore;
2601 use ironflow_store::models::StepStatus;
2602 use serde_json::json;
2603
2604 struct EchoWorkflow;
2606
2607 impl WorkflowHandler for EchoWorkflow {
2608 fn name(&self) -> &str {
2609 "echo-workflow"
2610 }
2611
2612 fn describe(&self) -> WorkflowInfo {
2613 WorkflowInfo {
2614 description: "A simple workflow that echoes hello".to_string(),
2615 source_code: None,
2616 sub_workflows: Vec::new(),
2617 category: None,
2618 version: self.version().map(str::to_string),
2619 compatible_versions: Vec::new(),
2620 input_schema: None,
2621 default_labels: HashMap::new(),
2622 schedule: self.schedule().cloned(),
2623 default_max_cost_usd: self.default_max_cost_usd(),
2624 }
2625 }
2626
2627 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2628 Box::pin(async move {
2629 ctx.shell("greet", ShellConfig::new("echo hello")).await?;
2630 Ok(())
2631 })
2632 }
2633 }
2634
2635 struct FailingWorkflow;
2637
2638 impl WorkflowHandler for FailingWorkflow {
2639 fn name(&self) -> &str {
2640 "failing-workflow"
2641 }
2642
2643 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2644 Box::pin(async move {
2645 ctx.shell("fail", ShellConfig::new("exit 1")).await?;
2646 Ok(())
2647 })
2648 }
2649 }
2650
2651 fn create_test_engine() -> Engine {
2652 let store = Arc::new(InMemoryStore::new());
2653 let inner = ClaudeCodeProvider::new();
2654 let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
2655 inner,
2656 "/tmp/ironflow-fixtures",
2657 ));
2658 Engine::new(store, provider)
2659 }
2660
2661 #[test]
2662 fn engine_new_creates_instance() {
2663 let engine = create_test_engine();
2664 assert_eq!(engine.handler_names().len(), 0);
2665 }
2666
2667 #[test]
2668 fn execution_mode_defaults_to_local() {
2669 let engine = create_test_engine();
2670 assert_eq!(engine.execution_mode(), ExecutionMode::Local);
2671 }
2672
2673 #[test]
2674 fn with_execution_mode_overrides_the_default() {
2675 let engine = create_test_engine().with_execution_mode(ExecutionMode::Workers);
2676 assert_eq!(engine.execution_mode(), ExecutionMode::Workers);
2677 }
2678
2679 #[test]
2680 fn engine_register_handler() {
2681 let mut engine = create_test_engine();
2682 let result = engine.register(EchoWorkflow);
2683 assert!(result.is_ok());
2684 assert_eq!(engine.handler_names().len(), 1);
2685 assert!(engine.handler_names().contains(&"echo-workflow"));
2686 }
2687
2688 #[test]
2689 fn engine_register_duplicate_returns_error() {
2690 let mut engine = create_test_engine();
2691 engine.register(EchoWorkflow).unwrap();
2692 let result = engine.register(EchoWorkflow);
2693 assert!(result.is_err());
2694 }
2695
2696 #[test]
2697 fn engine_get_handler_found() {
2698 let mut engine = create_test_engine();
2699 engine.register(EchoWorkflow).unwrap();
2700 let handler = engine.get_handler("echo-workflow");
2701 assert!(handler.is_some());
2702 }
2703
2704 #[test]
2705 fn engine_get_handler_not_found() {
2706 let engine = create_test_engine();
2707 let handler = engine.get_handler("nonexistent");
2708 assert!(handler.is_none());
2709 }
2710
2711 #[test]
2712 fn engine_handler_names_lists_all() {
2713 let mut engine = create_test_engine();
2714 engine.register(EchoWorkflow).unwrap();
2715 engine.register(FailingWorkflow).unwrap();
2716 let names = engine.handler_names();
2717 assert_eq!(names.len(), 2);
2718 assert!(names.contains(&"echo-workflow"));
2719 assert!(names.contains(&"failing-workflow"));
2720 }
2721
2722 #[test]
2723 fn engine_handler_info_returns_description() {
2724 let mut engine = create_test_engine();
2725 engine.register(EchoWorkflow).unwrap();
2726 let info = engine.handler_info("echo-workflow");
2727 assert!(info.is_some());
2728 let info = info.unwrap();
2729 assert_eq!(info.description, "A simple workflow that echoes hello");
2730 }
2731
2732 struct CategorizedWorkflow;
2733
2734 impl WorkflowHandler for CategorizedWorkflow {
2735 fn name(&self) -> &str {
2736 "categorized"
2737 }
2738 fn category(&self) -> Option<&str> {
2739 Some("data/etl")
2740 }
2741 fn execute<'a>(
2742 &'a self,
2743 _ctx: &'a mut WorkflowContext,
2744 ) -> crate::handler::HandlerFuture<'a> {
2745 Box::pin(async move { Ok(()) })
2746 }
2747 }
2748
2749 #[test]
2750 fn engine_default_describe_propagates_category() {
2751 let mut engine = create_test_engine();
2752 engine.register(CategorizedWorkflow).unwrap();
2753 let info = engine.handler_info("categorized").unwrap();
2754 assert_eq!(info.category.as_deref(), Some("data/etl"));
2755 }
2756
2757 #[test]
2758 fn engine_default_describe_without_category() {
2759 let mut engine = create_test_engine();
2760 engine.register(EchoWorkflow).unwrap();
2761 let info = engine.handler_info("echo-workflow").unwrap();
2762 assert!(info.category.is_none());
2763 }
2764
2765 struct ScheduledWorkflow {
2770 schedule: CronSchedule,
2771 }
2772
2773 impl ScheduledWorkflow {
2774 fn new() -> Self {
2775 Self {
2776 schedule: CronSchedule::new("0 0 * * * *").unwrap(),
2777 }
2778 }
2779 }
2780
2781 impl WorkflowHandler for ScheduledWorkflow {
2782 fn name(&self) -> &str {
2783 "scheduled"
2784 }
2785 fn schedule(&self) -> Option<&CronSchedule> {
2786 Some(&self.schedule)
2787 }
2788 fn execute<'a>(
2789 &'a self,
2790 _ctx: &'a mut WorkflowContext,
2791 ) -> crate::handler::HandlerFuture<'a> {
2792 Box::pin(async move { Ok(()) })
2793 }
2794 }
2795
2796 #[test]
2797 fn engine_default_describe_propagates_schedule() {
2798 let mut engine = create_test_engine();
2799 engine.register(ScheduledWorkflow::new()).unwrap();
2800 let info = engine.handler_info("scheduled").unwrap();
2801 assert_eq!(
2802 info.schedule.as_ref().map(|s| s.as_str()),
2803 Some("0 0 * * * *")
2804 );
2805 }
2806
2807 #[test]
2808 fn engine_default_describe_without_schedule() {
2809 let mut engine = create_test_engine();
2810 engine.register(EchoWorkflow).unwrap();
2811 let info = engine.handler_info("echo-workflow").unwrap();
2812 assert!(info.schedule.is_none());
2813 }
2814
2815 #[test]
2816 fn scheduled_handlers_returns_only_scheduled() {
2817 let mut engine = create_test_engine();
2818 engine.register(EchoWorkflow).unwrap();
2819 engine.register(ScheduledWorkflow::new()).unwrap();
2820 engine.register(FailingWorkflow).unwrap();
2821
2822 let scheduled = engine.scheduled_handlers();
2823 assert_eq!(scheduled.len(), 1);
2824 assert_eq!(scheduled[0].0, "scheduled");
2825 assert_eq!(scheduled[0].1.as_str(), "0 0 * * * *");
2826 }
2827
2828 #[test]
2829 fn scheduled_handlers_empty_when_none_scheduled() {
2830 let mut engine = create_test_engine();
2831 engine.register(EchoWorkflow).unwrap();
2832 engine.register(FailingWorkflow).unwrap();
2833
2834 let scheduled = engine.scheduled_handlers();
2835 assert!(scheduled.is_empty());
2836 }
2837
2838 struct BadCategoryWorkflow(&'static str);
2839
2840 impl WorkflowHandler for BadCategoryWorkflow {
2841 fn name(&self) -> &str {
2842 "bad-category"
2843 }
2844 fn category(&self) -> Option<&str> {
2845 Some(self.0)
2846 }
2847 fn execute<'a>(
2848 &'a self,
2849 _ctx: &'a mut WorkflowContext,
2850 ) -> crate::handler::HandlerFuture<'a> {
2851 Box::pin(async move { Ok(()) })
2852 }
2853 }
2854
2855 #[test]
2856 fn engine_register_rejects_empty_category() {
2857 let mut engine = create_test_engine();
2858 let err = engine.register(BadCategoryWorkflow("")).unwrap_err();
2859 match err {
2860 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty category")),
2861 other => panic!("expected InvalidWorkflow, got {other:?}"),
2862 }
2863 }
2864
2865 #[test]
2866 fn engine_register_rejects_leading_slash_category() {
2867 let mut engine = create_test_engine();
2868 let err = engine
2869 .register(BadCategoryWorkflow("/data/etl"))
2870 .unwrap_err();
2871 match err {
2872 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("leading '/'")),
2873 other => panic!("expected InvalidWorkflow, got {other:?}"),
2874 }
2875 }
2876
2877 #[test]
2878 fn engine_register_rejects_trailing_slash_category() {
2879 let mut engine = create_test_engine();
2880 let err = engine
2881 .register(BadCategoryWorkflow("data/etl/"))
2882 .unwrap_err();
2883 match err {
2884 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("trailing '/'")),
2885 other => panic!("expected InvalidWorkflow, got {other:?}"),
2886 }
2887 }
2888
2889 #[test]
2890 fn engine_register_rejects_double_slash_category() {
2891 let mut engine = create_test_engine();
2892 let err = engine
2893 .register(BadCategoryWorkflow("data//etl"))
2894 .unwrap_err();
2895 match err {
2896 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty segment")),
2897 other => panic!("expected InvalidWorkflow, got {other:?}"),
2898 }
2899 }
2900
2901 #[test]
2902 fn engine_register_rejects_whitespace_only_segment_category() {
2903 let mut engine = create_test_engine();
2904 let err = engine
2905 .register(BadCategoryWorkflow("data/ /etl"))
2906 .unwrap_err();
2907 match err {
2908 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("whitespace-only segment")),
2909 other => panic!("expected InvalidWorkflow, got {other:?}"),
2910 }
2911 }
2912
2913 #[test]
2914 fn engine_register_accepts_valid_nested_category() {
2915 let mut engine = create_test_engine();
2916 assert!(engine.register(CategorizedWorkflow).is_ok());
2917 }
2918
2919 #[tokio::test]
2920 async fn engine_unknown_workflow_returns_error() {
2921 let engine = create_test_engine();
2922 let result = engine
2923 .run_handler("unknown", TriggerKind::Manual, json!({}))
2924 .await;
2925 assert!(result.is_err());
2926 match result {
2927 Err(EngineError::InvalidWorkflow(msg)) => {
2928 assert!(msg.contains("no handler registered"));
2929 }
2930 _ => panic!("expected InvalidWorkflow error"),
2931 }
2932 }
2933
2934 #[tokio::test]
2935 async fn engine_enqueue_handler_creates_pending_run() {
2936 let mut engine = create_test_engine();
2937 engine.register(EchoWorkflow).unwrap();
2938
2939 let run = engine
2940 .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
2941 .await
2942 .unwrap();
2943 assert_eq!(run.status.state, RunStatus::Pending);
2944 assert_eq!(run.workflow_name, "echo-workflow");
2945 }
2946
2947 #[tokio::test]
2948 async fn enqueue_handler_leaves_the_run_unattributed() {
2949 let mut engine = create_test_engine();
2950 engine.register(EchoWorkflow).unwrap();
2951
2952 let run = engine
2953 .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
2954 .await
2955 .unwrap();
2956
2957 assert!(run.created_by.is_none());
2958 }
2959
2960 #[tokio::test]
2961 async fn enqueue_handler_with_options_records_the_author() {
2962 let mut engine = create_test_engine();
2963 engine.register(EchoWorkflow).unwrap();
2964 let actor = RunActor::User {
2965 user_id: Uuid::now_v7(),
2966 };
2967
2968 let run = engine
2969 .enqueue_handler_with_options(
2970 "echo-workflow",
2971 TriggerKind::Api,
2972 json!({}),
2973 EnqueueOptions {
2974 created_by: Some(actor.clone()),
2975 ..Default::default()
2976 },
2977 )
2978 .await
2979 .unwrap()
2980 .into_run();
2981
2982 assert_eq!(run.created_by, Some(actor));
2983 }
2984
2985 #[tokio::test]
2986 async fn enqueue_handler_with_options_accepts_no_author() {
2987 let mut engine = create_test_engine();
2988 engine.register(EchoWorkflow).unwrap();
2989
2990 let run = engine
2991 .enqueue_handler_with_options(
2992 "echo-workflow",
2993 TriggerKind::Cron {
2994 schedule: "0 * * * * *".to_string(),
2995 },
2996 json!({}),
2997 EnqueueOptions::default(),
2998 )
2999 .await
3000 .unwrap()
3001 .into_run();
3002
3003 assert!(run.created_by.is_none());
3004 }
3005
3006 #[tokio::test]
3007 async fn enqueue_handler_with_options_stores_concurrency_limits() {
3008 let mut engine = create_test_engine();
3009 engine.register(EchoWorkflow).unwrap();
3010 let limits = vec![
3011 ConcurrencyLimit::new("repo:acme", 2),
3012 ConcurrencyLimit::new("tenant:42", 5),
3013 ];
3014
3015 let run = engine
3016 .enqueue_handler_with_options(
3017 "echo-workflow",
3018 TriggerKind::Api,
3019 json!({}),
3020 EnqueueOptions {
3021 concurrency_limits: limits.clone(),
3022 ..Default::default()
3023 },
3024 )
3025 .await
3026 .unwrap()
3027 .into_run();
3028
3029 assert_eq!(run.concurrency_limits, limits);
3030 }
3031
3032 #[tokio::test]
3033 async fn enqueue_rejects_invalid_concurrency_limits() {
3034 let mut engine = create_test_engine();
3035 engine.register(EchoWorkflow).unwrap();
3036
3037 let invalid = [
3038 vec![ConcurrencyLimit::new("repo:acme", 0)],
3039 vec![ConcurrencyLimit::new("", 1)],
3040 vec![
3041 ConcurrencyLimit::new("repo:acme", 1),
3042 ConcurrencyLimit::new("repo:acme", 2),
3043 ],
3044 ];
3045 for concurrency_limits in invalid {
3046 let err = engine
3047 .enqueue_handler_with_options(
3048 "echo-workflow",
3049 TriggerKind::Api,
3050 json!({}),
3051 EnqueueOptions {
3052 concurrency_limits,
3053 ..Default::default()
3054 },
3055 )
3056 .await
3057 .unwrap_err();
3058 assert!(
3059 matches!(err, EngineError::InvalidConcurrencyLimit(_)),
3060 "{err:?}"
3061 );
3062 }
3063
3064 let err = engine
3066 .enqueue_handler_with_options(
3067 "not-registered",
3068 TriggerKind::Api,
3069 json!({}),
3070 EnqueueOptions {
3071 concurrency_limits: vec![ConcurrencyLimit::new("repo:acme", 0)],
3072 ..Default::default()
3073 },
3074 )
3075 .await
3076 .unwrap_err();
3077 assert!(
3078 matches!(err, EngineError::InvalidConcurrencyLimit(_)),
3079 "{err:?}"
3080 );
3081
3082 let page = engine
3083 .store()
3084 .list_runs(RunFilter::default(), 1, 10)
3085 .await
3086 .unwrap();
3087 assert_eq!(page.total, 0, "no run may be created");
3088 }
3089
3090 struct GpuWorkflow;
3091
3092 impl WorkflowHandler for GpuWorkflow {
3093 fn name(&self) -> &str {
3094 "gpu-workflow"
3095 }
3096
3097 fn required_worker_tags(&self) -> Vec<String> {
3098 vec!["gpu".to_string()]
3099 }
3100
3101 fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3102 Box::pin(async move { Ok(()) })
3103 }
3104 }
3105
3106 #[tokio::test]
3107 async fn enqueue_merges_handler_and_request_worker_tags() {
3108 let mut engine = create_test_engine();
3109 engine.register(GpuWorkflow).unwrap();
3110
3111 let run = engine
3112 .enqueue_handler_with_options(
3113 "gpu-workflow",
3114 TriggerKind::Api,
3115 json!({}),
3116 EnqueueOptions {
3117 worker_tags: vec!["region:eu".to_string(), "gpu".to_string()],
3118 ..Default::default()
3119 },
3120 )
3121 .await
3122 .unwrap()
3123 .into_run();
3124
3125 assert_eq!(
3126 run.worker_tags,
3127 vec!["gpu".to_string(), "region:eu".to_string()]
3128 );
3129 }
3130
3131 #[tokio::test]
3132 async fn enqueue_without_worker_tags_keeps_handler_tags() {
3133 let mut engine = create_test_engine();
3134 engine.register(GpuWorkflow).unwrap();
3135 engine.register(EchoWorkflow).unwrap();
3136
3137 let gpu = engine
3138 .enqueue_handler("gpu-workflow", TriggerKind::Api, json!({}), 0)
3139 .await
3140 .unwrap();
3141 assert_eq!(gpu.worker_tags, vec!["gpu".to_string()]);
3142
3143 let echo = engine
3144 .enqueue_handler("echo-workflow", TriggerKind::Api, json!({}), 0)
3145 .await
3146 .unwrap();
3147 assert!(echo.worker_tags.is_empty());
3148 }
3149
3150 #[tokio::test]
3151 async fn enqueue_rejects_invalid_worker_tags() {
3152 let mut engine = create_test_engine();
3153 engine.register(EchoWorkflow).unwrap();
3154
3155 for worker_tags in [
3156 vec!["bad,tag".to_string()],
3157 vec![" ".to_string()],
3158 vec!["x".repeat(65)],
3159 ] {
3160 let err = engine
3161 .enqueue_handler_with_options(
3162 "echo-workflow",
3163 TriggerKind::Api,
3164 json!({}),
3165 EnqueueOptions {
3166 worker_tags,
3167 ..Default::default()
3168 },
3169 )
3170 .await
3171 .unwrap_err();
3172 assert!(matches!(err, EngineError::InvalidWorkerTag(_)), "{err:?}");
3173 }
3174
3175 let err = engine
3177 .enqueue_handler_with_options(
3178 "not-registered",
3179 TriggerKind::Api,
3180 json!({}),
3181 EnqueueOptions {
3182 worker_tags: vec!["bad,tag".to_string()],
3183 ..Default::default()
3184 },
3185 )
3186 .await
3187 .unwrap_err();
3188 assert!(matches!(err, EngineError::InvalidWorkerTag(_)), "{err:?}");
3189 }
3190
3191 #[test]
3192 fn worker_tags_are_unset_by_default() {
3193 let engine = create_test_engine();
3194 assert!(engine.worker_tags().is_none());
3195 }
3196
3197 #[test]
3198 fn set_worker_tags_stores_the_tags() {
3199 let mut engine = create_test_engine();
3200 engine.set_worker_tags(vec!["gpu".to_string()]);
3201 assert_eq!(engine.worker_tags(), Some(&["gpu".to_string()][..]));
3202
3203 engine.set_worker_tags(Vec::new());
3204 assert_eq!(engine.worker_tags(), Some(&[][..]));
3205 }
3206
3207 #[tokio::test]
3208 async fn run_handler_records_handler_worker_tags() {
3209 let mut engine = create_test_engine();
3210 engine.register(GpuWorkflow).unwrap();
3211
3212 let result = engine
3213 .run_handler("gpu-workflow", TriggerKind::Manual, json!({}))
3214 .await
3215 .unwrap();
3216 assert_eq!(result.run.worker_tags, vec!["gpu".to_string()]);
3217 }
3218
3219 #[tokio::test]
3220 async fn run_handler_leaves_the_run_unattributed() {
3221 let mut engine = create_test_engine();
3222 engine.register(EchoWorkflow).unwrap();
3223
3224 let run = engine
3225 .run_handler("echo-workflow", TriggerKind::Manual, json!({}))
3226 .await
3227 .unwrap()
3228 .run;
3229
3230 assert!(run.created_by.is_none());
3231 }
3232
3233 #[tokio::test]
3234 async fn engine_register_boxed() {
3235 let mut engine = create_test_engine();
3236 let handler: Box<dyn WorkflowHandler> = Box::new(EchoWorkflow);
3237 let result = engine.register_boxed(handler);
3238 assert!(result.is_ok());
3239 assert_eq!(engine.handler_names().len(), 1);
3240 }
3241
3242 #[tokio::test]
3243 async fn engine_store_and_provider_accessors() {
3244 let store = Arc::new(InMemoryStore::new());
3245 let inner = ClaudeCodeProvider::new();
3246 let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
3247 inner,
3248 "/tmp/ironflow-fixtures",
3249 ));
3250 let engine = Engine::new(store.clone(), provider.clone());
3251
3252 let _ = engine.store();
3254 let _ = engine.provider();
3255 }
3256
3257 use crate::operation::{Operation, OperationContext};
3262 use async_trait::async_trait;
3263 use ironflow_core::error::OperationError;
3264 use ironflow_store::models::StepKind;
3265
3266 struct FakeGitlabOp {
3267 project_id: u64,
3268 title: String,
3269 }
3270
3271 #[async_trait]
3272 impl Operation for FakeGitlabOp {
3273 fn kind(&self) -> &str {
3274 "gitlab"
3275 }
3276
3277 async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
3278 Ok(json!({
3279 "issue_id": 42,
3280 "project_id": self.project_id,
3281 "title": self.title,
3282 }))
3283 }
3284
3285 fn input(&self) -> Option<Value> {
3286 Some(json!({
3287 "project_id": self.project_id,
3288 "title": self.title,
3289 }))
3290 }
3291 }
3292
3293 struct FailingOp;
3294
3295 #[async_trait]
3296 impl Operation for FailingOp {
3297 fn kind(&self) -> &str {
3298 "broken-service"
3299 }
3300
3301 async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
3302 Err(OperationError::Http {
3303 status: None,
3304 message: "service unavailable".to_string(),
3305 })
3306 }
3307 }
3308
3309 struct OperationWorkflow;
3310
3311 impl WorkflowHandler for OperationWorkflow {
3312 fn name(&self) -> &str {
3313 "operation-workflow"
3314 }
3315
3316 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3317 Box::pin(async move {
3318 let op = FakeGitlabOp {
3319 project_id: 123,
3320 title: "Bug report".to_string(),
3321 };
3322 ctx.operation("create-issue", &op).await?;
3323 Ok(())
3324 })
3325 }
3326 }
3327
3328 struct FailingOperationWorkflow;
3329
3330 impl WorkflowHandler for FailingOperationWorkflow {
3331 fn name(&self) -> &str {
3332 "failing-operation-workflow"
3333 }
3334
3335 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3336 Box::pin(async move {
3337 ctx.operation("broken-call", &FailingOp).await?;
3338 Ok(())
3339 })
3340 }
3341 }
3342
3343 struct MixedWorkflow;
3344
3345 impl WorkflowHandler for MixedWorkflow {
3346 fn name(&self) -> &str {
3347 "mixed-workflow"
3348 }
3349
3350 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3351 Box::pin(async move {
3352 ctx.shell("build", ShellConfig::new("echo built")).await?;
3353 let op = FakeGitlabOp {
3354 project_id: 456,
3355 title: "Deploy done".to_string(),
3356 };
3357 let result = ctx.operation("notify-gitlab", &op).await?;
3358 assert_eq!(result.output["issue_id"], 42);
3359 Ok(())
3360 })
3361 }
3362 }
3363
3364 #[tokio::test]
3365 async fn operation_step_happy_path() {
3366 let mut engine = create_test_engine();
3367 engine.register(OperationWorkflow).unwrap();
3368
3369 let run = engine
3370 .run_handler("operation-workflow", TriggerKind::Manual, json!({}))
3371 .await
3372 .unwrap()
3373 .run;
3374
3375 assert_eq!(run.status.state, RunStatus::Completed);
3376
3377 let steps = engine.store().list_steps(run.id).await.unwrap();
3378
3379 assert_eq!(steps.len(), 1);
3380 assert_eq!(steps[0].name, "create-issue");
3381 assert_eq!(steps[0].kind, StepKind::Custom("gitlab".to_string()));
3382 assert_eq!(
3383 steps[0].status.state,
3384 ironflow_store::models::StepStatus::Completed
3385 );
3386
3387 let output = steps[0].output.as_ref().unwrap();
3388 assert_eq!(output["issue_id"], 42);
3389 assert_eq!(output["project_id"], 123);
3390
3391 let input = steps[0].input.as_ref().unwrap();
3392 assert_eq!(input["project_id"], 123);
3393 assert_eq!(input["title"], "Bug report");
3394 }
3395
3396 #[tokio::test]
3397 async fn operation_step_failure_marks_run_failed() {
3398 let mut engine = create_test_engine();
3399 engine.register(FailingOperationWorkflow).unwrap();
3400
3401 let result = engine
3402 .run_handler("failing-operation-workflow", TriggerKind::Manual, json!({}))
3403 .await;
3404
3405 assert!(result.is_err());
3406 }
3407
3408 #[tokio::test]
3409 async fn operation_mixed_with_shell_steps() {
3410 let mut engine = create_test_engine();
3411 engine.register(MixedWorkflow).unwrap();
3412
3413 let run = engine
3414 .run_handler("mixed-workflow", TriggerKind::Manual, json!({}))
3415 .await
3416 .unwrap()
3417 .run;
3418
3419 assert_eq!(run.status.state, RunStatus::Completed);
3420
3421 let steps = engine.store().list_steps(run.id).await.unwrap();
3422
3423 assert_eq!(steps.len(), 2);
3424 assert_eq!(steps[0].kind, StepKind::Shell);
3425 assert_eq!(steps[1].kind, StepKind::Custom("gitlab".to_string()));
3426 assert_eq!(steps[0].position, 0);
3427 assert_eq!(steps[1].position, 1);
3428 }
3429
3430 use crate::config::ApprovalConfig;
3435
3436 struct SingleApprovalWorkflow;
3437
3438 impl WorkflowHandler for SingleApprovalWorkflow {
3439 fn name(&self) -> &str {
3440 "single-approval"
3441 }
3442
3443 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3444 Box::pin(async move {
3445 ctx.shell("build", ShellConfig::new("echo built")).await?;
3446 ctx.approval("gate", ApprovalConfig::new("OK?")).await?;
3447 ctx.shell("deploy", ShellConfig::new("echo deployed"))
3448 .await?;
3449 Ok(())
3450 })
3451 }
3452 }
3453
3454 struct DoubleApprovalWorkflow;
3455
3456 impl WorkflowHandler for DoubleApprovalWorkflow {
3457 fn name(&self) -> &str {
3458 "double-approval"
3459 }
3460
3461 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3462 Box::pin(async move {
3463 ctx.shell("build", ShellConfig::new("echo built")).await?;
3464 ctx.approval("staging-gate", ApprovalConfig::new("Deploy staging?"))
3465 .await?;
3466 ctx.shell("deploy-staging", ShellConfig::new("echo staging"))
3467 .await?;
3468 ctx.approval("prod-gate", ApprovalConfig::new("Deploy prod?"))
3469 .await?;
3470 ctx.shell("deploy-prod", ShellConfig::new("echo prod"))
3471 .await?;
3472 Ok(())
3473 })
3474 }
3475 }
3476
3477 #[tokio::test]
3478 async fn approval_pauses_run() {
3479 let mut engine = create_test_engine();
3480 engine.register(SingleApprovalWorkflow).unwrap();
3481
3482 let run = engine
3483 .run_handler("single-approval", TriggerKind::Manual, json!({}))
3484 .await
3485 .unwrap()
3486 .run;
3487
3488 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3489
3490 let steps = engine.store().list_steps(run.id).await.unwrap();
3491 assert_eq!(steps.len(), 2); assert_eq!(steps[0].kind, StepKind::Shell);
3493 assert_eq!(steps[0].status.state, StepStatus::Completed);
3494 assert_eq!(steps[1].kind, StepKind::Approval);
3495 assert_eq!(steps[1].status.state, StepStatus::AwaitingApproval);
3496 }
3497
3498 #[tokio::test]
3499 async fn approval_resume_completes_run() {
3500 let mut engine = create_test_engine();
3501 engine.register(SingleApprovalWorkflow).unwrap();
3502
3503 let run = engine
3505 .run_handler("single-approval", TriggerKind::Manual, json!({}))
3506 .await
3507 .unwrap()
3508 .run;
3509 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3510
3511 engine
3513 .store()
3514 .update_run_status(run.id, RunStatus::Running)
3515 .await
3516 .unwrap();
3517
3518 let resumed = engine.resume_run(run.id).await.unwrap().run;
3520 assert_eq!(resumed.status.state, RunStatus::Completed);
3521
3522 let steps = engine.store().list_steps(run.id).await.unwrap();
3523 assert_eq!(steps.len(), 3); assert_eq!(steps[0].name, "build");
3525 assert_eq!(steps[0].status.state, StepStatus::Completed);
3526 assert_eq!(steps[1].name, "gate");
3527 assert_eq!(steps[1].kind, StepKind::Approval);
3528 assert_eq!(steps[1].status.state, StepStatus::Completed);
3529 assert_eq!(steps[2].name, "deploy");
3530 assert_eq!(steps[2].status.state, StepStatus::Completed);
3531 }
3532
3533 #[tokio::test]
3534 async fn double_approval_two_resumes() {
3535 let mut engine = create_test_engine();
3536 engine.register(DoubleApprovalWorkflow).unwrap();
3537
3538 let run = engine
3540 .run_handler("double-approval", TriggerKind::Manual, json!({}))
3541 .await
3542 .unwrap()
3543 .run;
3544 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3545
3546 let steps = engine.store().list_steps(run.id).await.unwrap();
3547 assert_eq!(steps.len(), 2); engine
3551 .store()
3552 .update_run_status(run.id, RunStatus::Running)
3553 .await
3554 .unwrap();
3555
3556 let resumed = engine.resume_run(run.id).await.unwrap().run;
3557 assert_eq!(resumed.status.state, RunStatus::AwaitingApproval);
3558
3559 let steps = engine.store().list_steps(run.id).await.unwrap();
3560 assert_eq!(steps.len(), 4); engine
3564 .store()
3565 .update_run_status(run.id, RunStatus::Running)
3566 .await
3567 .unwrap();
3568
3569 let final_run = engine.resume_run(run.id).await.unwrap().run;
3570 assert_eq!(final_run.status.state, RunStatus::Completed);
3571
3572 let steps = engine.store().list_steps(run.id).await.unwrap();
3573 assert_eq!(steps.len(), 5);
3574 assert_eq!(steps[0].name, "build");
3575 assert_eq!(steps[1].name, "staging-gate");
3576 assert_eq!(steps[2].name, "deploy-staging");
3577 assert_eq!(steps[3].name, "prod-gate");
3578 assert_eq!(steps[4].name, "deploy-prod");
3579
3580 for step in &steps {
3581 assert_eq!(step.status.state, StepStatus::Completed);
3582 }
3583 }
3584
3585 use ironflow_store::models::{NewStep, StepUpdate, step_trace_id};
3590
3591 async fn create_step_with_status(
3592 store: &Arc<dyn Store>,
3593 run_id: Uuid,
3594 name: &str,
3595 position: u32,
3596 status: StepStatus,
3597 ) -> ironflow_store::models::Step {
3598 let step = store
3599 .create_step(NewStep {
3600 run_id,
3601 trace_id: step_trace_id(run_id, name, position),
3602 name: name.to_string(),
3603 kind: StepKind::Shell,
3604 position,
3605 input: None,
3606 is_error_handler: false,
3607 })
3608 .await
3609 .unwrap();
3610
3611 match status {
3612 StepStatus::Pending => {}
3613 StepStatus::Running => {
3614 store
3615 .update_step(
3616 step.id,
3617 StepUpdate {
3618 status: Some(StepStatus::Running),
3619 ..StepUpdate::default()
3620 },
3621 )
3622 .await
3623 .unwrap();
3624 }
3625 StepStatus::Completed => {
3626 store
3627 .update_step(
3628 step.id,
3629 StepUpdate {
3630 status: Some(StepStatus::Running),
3631 ..StepUpdate::default()
3632 },
3633 )
3634 .await
3635 .unwrap();
3636 store
3637 .update_step(
3638 step.id,
3639 StepUpdate {
3640 status: Some(StepStatus::Completed),
3641 ..StepUpdate::default()
3642 },
3643 )
3644 .await
3645 .unwrap();
3646 }
3647 StepStatus::AwaitingApproval => {
3648 store
3649 .update_step(
3650 step.id,
3651 StepUpdate {
3652 status: Some(StepStatus::Running),
3653 ..StepUpdate::default()
3654 },
3655 )
3656 .await
3657 .unwrap();
3658 store
3659 .update_step(
3660 step.id,
3661 StepUpdate {
3662 status: Some(StepStatus::AwaitingApproval),
3663 ..StepUpdate::default()
3664 },
3665 )
3666 .await
3667 .unwrap();
3668 }
3669 _ => panic!("unsupported status for test helper: {status}"),
3670 }
3671
3672 store.get_step(step.id).await.unwrap().unwrap()
3673 }
3674
3675 #[tokio::test]
3676 async fn fail_orphaned_steps_marks_running_as_failed() {
3677 let engine = create_test_engine();
3678 let run = engine
3679 .store()
3680 .create_run(NewRun {
3681 created_by: None,
3682 workflow_name: "test".to_string(),
3683 trigger: TriggerKind::Manual,
3684 payload: json!({}),
3685 max_retries: 0,
3686 handler_version: None,
3687 labels: HashMap::new(),
3688 scheduled_at: None,
3689 idempotency_key: None,
3690 concurrency_key: None,
3691 concurrency_limits: Vec::new(),
3692 max_cost_usd: None,
3693 worker_tags: Vec::new(),
3694 })
3695 .await
3696 .unwrap()
3697 .into_run();
3698
3699 let step = create_step_with_status(
3700 engine.store(),
3701 run.id,
3702 "running-step",
3703 0,
3704 StepStatus::Running,
3705 )
3706 .await;
3707
3708 engine
3709 .fail_orphaned_steps(run.id, "parent run timed out")
3710 .await
3711 .unwrap();
3712
3713 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
3714 assert_eq!(updated.status.state, StepStatus::Failed);
3715 assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
3716 assert!(updated.completed_at.is_some());
3717 }
3718
3719 #[tokio::test]
3720 async fn fail_orphaned_steps_marks_pending_as_skipped() {
3721 let engine = create_test_engine();
3722 let run = engine
3723 .store()
3724 .create_run(NewRun {
3725 created_by: None,
3726 workflow_name: "test".to_string(),
3727 trigger: TriggerKind::Manual,
3728 payload: json!({}),
3729 max_retries: 0,
3730 handler_version: None,
3731 labels: HashMap::new(),
3732 scheduled_at: None,
3733 idempotency_key: None,
3734 concurrency_key: None,
3735 concurrency_limits: Vec::new(),
3736 max_cost_usd: None,
3737 worker_tags: Vec::new(),
3738 })
3739 .await
3740 .unwrap()
3741 .into_run();
3742
3743 let step = create_step_with_status(
3744 engine.store(),
3745 run.id,
3746 "pending-step",
3747 0,
3748 StepStatus::Pending,
3749 )
3750 .await;
3751
3752 engine
3753 .fail_orphaned_steps(run.id, "parent run timed out")
3754 .await
3755 .unwrap();
3756
3757 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
3758 assert_eq!(updated.status.state, StepStatus::Skipped);
3759 assert!(updated.error.is_none());
3760 assert!(updated.completed_at.is_some());
3761 }
3762
3763 #[tokio::test]
3764 async fn fail_orphaned_steps_marks_awaiting_approval_as_failed() {
3765 let engine = create_test_engine();
3766 let run = engine
3767 .store()
3768 .create_run(NewRun {
3769 created_by: None,
3770 workflow_name: "test".to_string(),
3771 trigger: TriggerKind::Manual,
3772 payload: json!({}),
3773 max_retries: 0,
3774 handler_version: None,
3775 labels: HashMap::new(),
3776 scheduled_at: None,
3777 idempotency_key: None,
3778 concurrency_key: None,
3779 concurrency_limits: Vec::new(),
3780 max_cost_usd: None,
3781 worker_tags: Vec::new(),
3782 })
3783 .await
3784 .unwrap()
3785 .into_run();
3786
3787 let step = create_step_with_status(
3788 engine.store(),
3789 run.id,
3790 "approval-step",
3791 0,
3792 StepStatus::AwaitingApproval,
3793 )
3794 .await;
3795
3796 engine
3797 .fail_orphaned_steps(run.id, "parent run timed out")
3798 .await
3799 .unwrap();
3800
3801 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
3802 assert_eq!(updated.status.state, StepStatus::Failed);
3803 assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
3804 assert!(updated.completed_at.is_some());
3805 }
3806
3807 #[tokio::test]
3808 async fn fail_orphaned_steps_skips_terminal_steps() {
3809 let engine = create_test_engine();
3810 let run = engine
3811 .store()
3812 .create_run(NewRun {
3813 created_by: None,
3814 workflow_name: "test".to_string(),
3815 trigger: TriggerKind::Manual,
3816 payload: json!({}),
3817 max_retries: 0,
3818 handler_version: None,
3819 labels: HashMap::new(),
3820 scheduled_at: None,
3821 idempotency_key: None,
3822 concurrency_key: None,
3823 concurrency_limits: Vec::new(),
3824 max_cost_usd: None,
3825 worker_tags: Vec::new(),
3826 })
3827 .await
3828 .unwrap()
3829 .into_run();
3830
3831 let completed_step =
3832 create_step_with_status(engine.store(), run.id, "done", 0, StepStatus::Completed).await;
3833 let running_step =
3834 create_step_with_status(engine.store(), run.id, "in-flight", 1, StepStatus::Running)
3835 .await;
3836
3837 engine
3838 .fail_orphaned_steps(run.id, "parent run timed out")
3839 .await
3840 .unwrap();
3841
3842 let completed = engine
3843 .store()
3844 .get_step(completed_step.id)
3845 .await
3846 .unwrap()
3847 .unwrap();
3848 assert_eq!(completed.status.state, StepStatus::Completed);
3849
3850 let failed = engine
3851 .store()
3852 .get_step(running_step.id)
3853 .await
3854 .unwrap()
3855 .unwrap();
3856 assert_eq!(failed.status.state, StepStatus::Failed);
3857 }
3858
3859 #[tokio::test]
3860 async fn fail_orphaned_steps_mixed_states() {
3861 let engine = create_test_engine();
3862 let run = engine
3863 .store()
3864 .create_run(NewRun {
3865 created_by: None,
3866 workflow_name: "test".to_string(),
3867 trigger: TriggerKind::Manual,
3868 payload: json!({}),
3869 max_retries: 0,
3870 handler_version: None,
3871 labels: HashMap::new(),
3872 scheduled_at: None,
3873 idempotency_key: None,
3874 concurrency_key: None,
3875 concurrency_limits: Vec::new(),
3876 max_cost_usd: None,
3877 worker_tags: Vec::new(),
3878 })
3879 .await
3880 .unwrap()
3881 .into_run();
3882
3883 let s_completed =
3884 create_step_with_status(engine.store(), run.id, "step-1", 0, StepStatus::Completed)
3885 .await;
3886 let s_running =
3887 create_step_with_status(engine.store(), run.id, "step-2", 1, StepStatus::Running).await;
3888 let s_pending =
3889 create_step_with_status(engine.store(), run.id, "step-3", 2, StepStatus::Pending).await;
3890
3891 engine.fail_orphaned_steps(run.id, "timeout").await.unwrap();
3892
3893 let r_completed = engine
3894 .store()
3895 .get_step(s_completed.id)
3896 .await
3897 .unwrap()
3898 .unwrap();
3899 assert_eq!(r_completed.status.state, StepStatus::Completed);
3900
3901 let r_running = engine
3902 .store()
3903 .get_step(s_running.id)
3904 .await
3905 .unwrap()
3906 .unwrap();
3907 assert_eq!(r_running.status.state, StepStatus::Failed);
3908 assert_eq!(r_running.error.as_deref(), Some("timeout"));
3909
3910 let r_pending = engine
3911 .store()
3912 .get_step(s_pending.id)
3913 .await
3914 .unwrap()
3915 .unwrap();
3916 assert_eq!(r_pending.status.state, StepStatus::Skipped);
3917 assert!(r_pending.error.is_none());
3918 }
3919
3920 #[tokio::test]
3921 async fn fail_orphaned_steps_no_steps_is_noop() {
3922 let engine = create_test_engine();
3923 let run = engine
3924 .store()
3925 .create_run(NewRun {
3926 created_by: None,
3927 workflow_name: "test".to_string(),
3928 trigger: TriggerKind::Manual,
3929 payload: json!({}),
3930 max_retries: 0,
3931 handler_version: None,
3932 labels: HashMap::new(),
3933 scheduled_at: None,
3934 idempotency_key: None,
3935 concurrency_key: None,
3936 concurrency_limits: Vec::new(),
3937 max_cost_usd: None,
3938 worker_tags: Vec::new(),
3939 })
3940 .await
3941 .unwrap()
3942 .into_run();
3943
3944 let result = engine.fail_orphaned_steps(run.id, "timeout").await;
3945 assert!(result.is_ok());
3946 }
3947
3948 #[tokio::test]
3949 async fn fail_orphaned_steps_preserves_existing_error() {
3950 let engine = create_test_engine();
3951 let run = engine
3952 .store()
3953 .create_run(NewRun {
3954 created_by: None,
3955 workflow_name: "test".to_string(),
3956 trigger: TriggerKind::Manual,
3957 payload: json!({}),
3958 max_retries: 0,
3959 handler_version: None,
3960 labels: HashMap::new(),
3961 scheduled_at: None,
3962 idempotency_key: None,
3963 concurrency_key: None,
3964 concurrency_limits: Vec::new(),
3965 max_cost_usd: None,
3966 worker_tags: Vec::new(),
3967 })
3968 .await
3969 .unwrap()
3970 .into_run();
3971
3972 let step_with_error = create_step_with_status(
3973 engine.store(),
3974 run.id,
3975 "already-errored",
3976 0,
3977 StepStatus::Running,
3978 )
3979 .await;
3980
3981 engine
3982 .store()
3983 .update_step(
3984 step_with_error.id,
3985 StepUpdate {
3986 error: Some("real error from provider".to_string()),
3987 ..StepUpdate::default()
3988 },
3989 )
3990 .await
3991 .unwrap();
3992
3993 let step_no_error = create_step_with_status(
3994 engine.store(),
3995 run.id,
3996 "no-error-yet",
3997 1,
3998 StepStatus::Running,
3999 )
4000 .await;
4001
4002 engine
4003 .fail_orphaned_steps(run.id, "parent run failed")
4004 .await
4005 .unwrap();
4006
4007 let updated_with = engine
4008 .store()
4009 .get_step(step_with_error.id)
4010 .await
4011 .unwrap()
4012 .unwrap();
4013 assert_eq!(updated_with.status.state, StepStatus::Failed);
4014 assert_eq!(
4015 updated_with.error.as_deref(),
4016 Some("real error from provider"),
4017 );
4018
4019 let updated_without = engine
4020 .store()
4021 .get_step(step_no_error.id)
4022 .await
4023 .unwrap()
4024 .unwrap();
4025 assert_eq!(updated_without.status.state, StepStatus::Failed);
4026 assert_eq!(updated_without.error.as_deref(), Some("parent run failed"),);
4027 }
4028}