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, validate_concurrency_limits,
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}
141
142#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
153pub enum ExecutionMode {
154 #[default]
158 Local,
159 Workers,
162}
163
164pub struct Engine {
205 store: Arc<dyn Store>,
206 provider: Arc<dyn AgentProvider>,
207 handlers: HashMap<String, Arc<dyn WorkflowHandler>>,
208 event_publisher: EventPublisher,
209 log_sender: Option<LogSender>,
210 budget: BudgetConfig,
211 artifact_sink: Option<Arc<dyn ArtifactSink>>,
212 guard_config: Option<WorkflowGuardConfig>,
213 event_bus: Option<WorkflowEventBus>,
214 decision_provider: Option<Arc<dyn DecisionProvider>>,
215 step_interceptor: Option<Arc<dyn StepInterceptor>>,
216 execution_mode: ExecutionMode,
217}
218
219fn validate_category(handler_name: &str, category: &str) -> Result<(), EngineError> {
229 let reject = |reason: &str| {
230 Err(EngineError::InvalidWorkflow(format!(
231 "handler '{handler_name}' has invalid category '{category}': {reason}"
232 )))
233 };
234
235 if category.is_empty() {
236 return reject("empty category");
237 }
238 if category.starts_with('/') {
239 return reject("leading '/'");
240 }
241 if category.ends_with('/') {
242 return reject("trailing '/'");
243 }
244 for segment in category.split('/') {
245 if segment.is_empty() {
246 return reject("empty segment (double '/')");
247 }
248 if segment.trim().is_empty() {
249 return reject("whitespace-only segment");
250 }
251 }
252 Ok(())
253}
254
255fn chain_label(run: &Run, key: &str) -> Option<Uuid> {
260 if !matches!(run.trigger, TriggerKind::Workflow) {
261 return None;
262 }
263 let id = Uuid::parse_str(run.labels.get(key)?).ok()?;
264 (id != run.id).then_some(id)
265}
266
267pub fn chain_root(run: &Run) -> Option<Uuid> {
287 chain_label(run, LABEL_ROOT_RUN_ID)
288}
289
290fn chain_parent(run: &Run) -> Option<Uuid> {
292 chain_label(run, PARENT_RUN_ID_LABEL)
293}
294
295impl Engine {
296 pub fn new(store: Arc<dyn Store>, provider: Arc<dyn AgentProvider>) -> Self {
312 Self {
313 store,
314 provider,
315 handlers: HashMap::new(),
316 event_publisher: EventPublisher::new(),
317 log_sender: None,
318 budget: BudgetConfig::new(),
319 artifact_sink: None,
320 guard_config: None,
321 event_bus: None,
322 decision_provider: None,
323 step_interceptor: None,
324 execution_mode: ExecutionMode::default(),
325 }
326 }
327
328 pub fn with_decision_provider(mut self, provider: Arc<dyn DecisionProvider>) -> Self {
349 self.decision_provider = Some(provider);
350 self
351 }
352
353 pub fn with_step_interceptor(mut self, interceptor: Arc<dyn StepInterceptor>) -> Self {
377 self.step_interceptor = Some(interceptor);
378 self
379 }
380
381 pub fn step_interceptor(&self) -> Option<&Arc<dyn StepInterceptor>> {
383 self.step_interceptor.as_ref()
384 }
385
386 pub fn with_budget_config(mut self, budget: BudgetConfig) -> Self {
407 self.budget = budget;
408 self
409 }
410
411 pub fn budget_config(&self) -> &BudgetConfig {
413 &self.budget
414 }
415
416 pub fn with_guard_config(mut self, config: WorkflowGuardConfig) -> Self {
438 self.guard_config = Some(config);
439 self
440 }
441
442 pub fn guard_config(&self) -> Option<&WorkflowGuardConfig> {
444 self.guard_config.as_ref()
445 }
446
447 pub fn with_execution_mode(mut self, mode: ExecutionMode) -> Self {
468 self.execution_mode = mode;
469 self
470 }
471
472 pub fn execution_mode(&self) -> ExecutionMode {
474 self.execution_mode
475 }
476
477 pub fn set_log_sender(&mut self, sender: LogSender) {
483 self.log_sender = Some(sender);
484 }
485
486 pub fn set_artifact_sink(&mut self, sink: Arc<dyn ArtifactSink>) {
504 self.artifact_sink = Some(sink);
505 }
506
507 pub fn artifact_sink(&self) -> Option<&Arc<dyn ArtifactSink>> {
509 self.artifact_sink.as_ref()
510 }
511
512 pub fn set_event_bus(&mut self, bus: WorkflowEventBus) {
529 self.event_bus = Some(bus);
530 }
531
532 pub fn event_bus(&self) -> Option<&WorkflowEventBus> {
534 self.event_bus.as_ref()
535 }
536
537 pub fn store(&self) -> &Arc<dyn Store> {
539 &self.store
540 }
541
542 pub fn provider(&self) -> &Arc<dyn AgentProvider> {
544 &self.provider
545 }
546
547 fn build_context(&self, run: &Run) -> WorkflowContext {
556 let handlers = self.handlers.clone();
557 let resolver: crate::context::HandlerResolver =
558 Arc::new(move |name: &str| handlers.get(name).cloned());
559 let mut ctx = WorkflowContext::with_handler_resolver(
560 run.id,
561 run.workflow_name.clone(),
562 self.store.clone(),
563 self.provider.clone(),
564 resolver,
565 );
566 ctx.carry_over_run_totals(run.retry_count + 1, run.cost_usd, run.duration_ms);
567 ctx.set_max_cost_usd(run.max_cost_usd);
568 ctx.set_run_created_at(run.created_at);
569 if let Some(ref sender) = self.log_sender {
570 ctx.set_log_sender(sender.clone());
571 }
572 if let Some(ref sink) = self.artifact_sink {
573 ctx.set_artifact_sink(sink.clone());
574 }
575 if let Some(ref bus) = self.event_bus {
576 ctx.set_event_bus(bus.clone());
577 }
578 if let Some(ref provider) = self.decision_provider {
579 ctx.set_decision_provider(provider.clone());
580 }
581 if let Some(ref interceptor) = self.step_interceptor {
582 ctx.set_step_interceptor(interceptor.clone());
583 }
584 ctx
585 }
586
587 fn build_context_with_guard(
593 &self,
594 run: &Run,
595 handler: &dyn WorkflowHandler,
596 ) -> WorkflowContext {
597 let mut ctx = self.build_context(run);
598 let guard_config = handler.guard_config().or_else(|| self.guard_config.clone());
599 if let Some(config) = guard_config {
600 ctx.set_guard(config, new_shared_guard_state());
601 }
602 ctx
603 }
604
605 async fn check_monthly_quota(&self, workflow_name: &str) -> Result<(), EngineError> {
616 let Some(limit) = self.budget.monthly_cost_limit_usd else {
617 return Ok(());
618 };
619
620 let stats = self
621 .store
622 .get_stats(RunFilter {
623 created_after: Some(month_start(Utc::now())),
624 ..RunFilter::default()
625 })
626 .await?;
627
628 if stats.total_cost_usd < limit {
629 return Ok(());
630 }
631
632 warn!(
633 workflow = %workflow_name,
634 limit_usd = %limit,
635 spent_usd = %stats.total_cost_usd,
636 "monthly cost quota exhausted, refusing new run"
637 );
638
639 #[cfg(feature = "prometheus")]
640 counter!(
641 RUN_BUDGET_EXCEEDED_TOTAL,
642 "workflow" => workflow_name.to_string(),
643 "scope" => "monthly",
644 )
645 .increment(1);
646
647 Err(EngineError::MonthlyBudgetExceeded {
648 limit_usd: limit,
649 spent_usd: stats.total_cost_usd,
650 })
651 }
652
653 pub fn register(&mut self, handler: impl WorkflowHandler + 'static) -> Result<(), EngineError> {
697 let name = handler.name().to_string();
698 if self.handlers.contains_key(&name) {
699 return Err(EngineError::InvalidWorkflow(format!(
700 "handler '{}' already registered",
701 name
702 )));
703 }
704 if let Some(category) = handler.category() {
705 validate_category(&name, category)?;
706 }
707 self.handlers.insert(name, Arc::new(handler));
708 Ok(())
709 }
710
711 pub fn register_boxed(&mut self, handler: Box<dyn WorkflowHandler>) -> Result<(), EngineError> {
718 let name = handler.name().to_string();
719 if self.handlers.contains_key(&name) {
720 return Err(EngineError::InvalidWorkflow(format!(
721 "handler '{}' already registered",
722 name
723 )));
724 }
725 if let Some(category) = handler.category() {
726 validate_category(&name, category)?;
727 }
728 self.handlers.insert(name, Arc::from(handler));
729 Ok(())
730 }
731
732 pub fn get_handler(&self, name: &str) -> Option<&Arc<dyn WorkflowHandler>> {
734 self.handlers.get(name)
735 }
736
737 pub fn handler_names(&self) -> Vec<&str> {
739 self.handlers.keys().map(|s| s.as_str()).collect()
740 }
741
742 pub fn handler_info(&self, name: &str) -> Option<WorkflowInfo> {
744 self.handlers.get(name).map(|h| h.describe())
745 }
746
747 pub fn scheduled_handlers(&self) -> Vec<(&str, &CronSchedule)> {
771 self.handlers
772 .iter()
773 .filter_map(|(name, handler)| handler.schedule().map(|sched| (name.as_str(), sched)))
774 .collect()
775 }
776
777 pub fn subscribe(
802 &mut self,
803 subscriber: impl EventSubscriber + 'static,
804 event_types: &[&'static str],
805 ) {
806 self.event_publisher.subscribe(subscriber, event_types);
807 }
808
809 pub fn event_publisher(&self) -> &EventPublisher {
814 &self.event_publisher
815 }
816
817 #[tracing::instrument(name = "engine.run_handler", skip_all, fields(workflow = %handler_name))]
847 pub async fn run_handler(
848 &self,
849 handler_name: &str,
850 trigger: TriggerKind,
851 payload: Value,
852 ) -> Result<WorkflowResult, EngineError> {
853 let handler = self
854 .handlers
855 .get(handler_name)
856 .ok_or_else(|| {
857 EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
858 })?
859 .clone();
860
861 self.check_monthly_quota(handler_name).await?;
862
863 let handler_version = handler.version().map(str::to_string);
864 let max_cost_usd = self
865 .budget
866 .resolve_run_cap(None, handler.default_max_cost_usd());
867 let run = self
868 .store
869 .create_run(NewRun {
870 created_by: None,
871 workflow_name: handler_name.to_string(),
872 trigger,
873 payload,
874 max_retries: 0,
875 handler_version,
876 labels: handler.default_labels(),
877 scheduled_at: None,
878 idempotency_key: None,
879 concurrency_key: None,
880 concurrency_limits: Vec::new(),
881 max_cost_usd,
882 })
883 .await?
884 .into_run();
885
886 let run_id = run.id;
887 info!(run_id = %run_id, handler_version = run.handler_version.as_deref().unwrap_or(""), "run created");
888
889 self.store
890 .update_run_status(run_id, RunStatus::Running)
891 .await?;
892
893 #[cfg(feature = "prometheus")]
894 gauge!(RUNS_ACTIVE, "workflow" => handler_name.to_string()).increment(1.0);
895
896 let run_start = Instant::now();
897 let mut ctx = self.build_context_with_guard(&run, handler.as_ref());
898
899 let result = handler.execute(&mut ctx).await;
900 self.finalize_run(run_id, handler_name, result, &ctx, run_start, run.labels)
901 .await
902 }
903
904 #[tracing::instrument(name = "engine.plan_handler", skip_all, fields(workflow = %handler_name))]
946 pub async fn plan_handler(
947 &self,
948 handler_name: &str,
949 payload: Value,
950 options: PlanOptions,
951 ) -> Result<ExecutionPlan, EngineError> {
952 if options.max_depth == 0 {
953 return Err(EngineError::InvalidWorkflow(
954 "max_depth must be at least 1".to_string(),
955 ));
956 }
957
958 let handler = self
959 .handlers
960 .get(handler_name)
961 .ok_or_else(|| {
962 EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
963 })?
964 .clone();
965
966 let estimates = if options.estimate_durations {
967 estimate_durations(&self.store, handler_name, options.sample_runs).await?
968 } else {
969 HashMap::new()
970 };
971
972 let shared: SharedPlanRecorder = Arc::new(Mutex::new(PlanRecorder::new(
973 handler_name.to_string(),
974 payload,
975 options.max_depth,
976 estimates,
977 )));
978
979 let handlers = self.handlers.clone();
982 let resolver: crate::context::HandlerResolver =
983 Arc::new(move |name: &str| handlers.get(name).cloned());
984 let mut ctx = WorkflowContext::with_handler_resolver(
985 Uuid::now_v7(),
986 handler_name.to_string(),
987 self.store.clone(),
988 self.provider.clone(),
989 resolver,
990 );
991 ctx.set_plan(shared.clone());
992
993 if let Err(err) = handler.execute(&mut ctx).await {
994 lock_plan(&shared).fail(err.to_string());
995 }
996 drop(ctx);
997
998 let plan = match Arc::try_unwrap(shared) {
999 Ok(mutex) => mutex
1000 .into_inner()
1001 .unwrap_or_else(|poisoned| poisoned.into_inner())
1002 .into_plan(),
1003 Err(shared) => lock_plan(&shared).snapshot(),
1004 };
1005
1006 info!(
1007 workflow = %handler_name,
1008 steps = plan.steps.len(),
1009 truncated = plan.truncated,
1010 "execution plan built"
1011 );
1012
1013 Ok(plan)
1014 }
1015
1016 #[tracing::instrument(name = "engine.enqueue_handler", skip_all, fields(workflow = %handler_name))]
1027 pub async fn enqueue_handler(
1028 &self,
1029 handler_name: &str,
1030 trigger: TriggerKind,
1031 payload: Value,
1032 max_retries: u32,
1033 ) -> Result<Run, EngineError> {
1034 self.enqueue_handler_with_options(
1035 handler_name,
1036 trigger,
1037 payload,
1038 EnqueueOptions {
1039 max_retries,
1040 ..Default::default()
1041 },
1042 )
1043 .await
1044 .map(RunCreation::into_run)
1045 }
1046
1047 #[tracing::instrument(name = "engine.enqueue_handler_with_options", skip_all, fields(workflow = %handler_name))]
1095 pub async fn enqueue_handler_with_options(
1096 &self,
1097 handler_name: &str,
1098 trigger: TriggerKind,
1099 payload: Value,
1100 options: EnqueueOptions,
1101 ) -> Result<RunCreation, EngineError> {
1102 let EnqueueOptions {
1103 max_retries,
1104 labels,
1105 scheduled_at,
1106 max_cost_usd,
1107 created_by,
1108 idempotency_key,
1109 concurrency_key,
1110 concurrency_limits,
1111 } = options;
1112
1113 validate_concurrency_limits(&concurrency_limits)
1116 .map_err(EngineError::InvalidConcurrencyLimit)?;
1117
1118 let handler = self.handlers.get(handler_name).ok_or_else(|| {
1119 EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
1120 })?;
1121
1122 self.check_monthly_quota(handler_name).await?;
1123
1124 let handler_version = handler.version().map(str::to_string);
1125 let mut merged_labels = handler.default_labels();
1126 merged_labels.extend(labels);
1127 let resolved_cap = self
1128 .budget
1129 .resolve_run_cap(max_cost_usd, handler.default_max_cost_usd());
1130
1131 let creation = self
1132 .store
1133 .create_run(NewRun {
1134 workflow_name: handler_name.to_string(),
1135 trigger,
1136 payload,
1137 max_retries,
1138 handler_version,
1139 labels: merged_labels,
1140 scheduled_at,
1141 created_by,
1142 idempotency_key,
1143 concurrency_key,
1144 concurrency_limits,
1145 max_cost_usd: resolved_cap,
1146 })
1147 .await?;
1148
1149 match &creation {
1150 RunCreation::Created(run) => info!(
1151 run_id = %run.id,
1152 workflow = %handler_name,
1153 max_cost_usd = ?resolved_cap,
1154 "handler run enqueued"
1155 ),
1156 RunCreation::Existing(run) => info!(
1157 run_id = %run.id,
1158 workflow = %handler_name,
1159 "idempotent replay, nothing enqueued"
1160 ),
1161 }
1162
1163 Ok(creation)
1164 }
1165
1166 #[tracing::instrument(name = "engine.execute_handler_run", skip_all, fields(run_id = %run_id))]
1186 pub async fn execute_handler_run(&self, run_id: Uuid) -> Result<WorkflowResult, EngineError> {
1187 let run = self
1188 .store
1189 .get_run(run_id)
1190 .await?
1191 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
1192
1193 if let Some(root_run_id) = chain_root(&run) {
1194 return self.resume_chain(run, root_run_id).await;
1195 }
1196
1197 let handler = self
1198 .handlers
1199 .get(&run.workflow_name)
1200 .ok_or_else(|| {
1201 EngineError::InvalidWorkflow(format!(
1202 "no handler registered: {}",
1203 run.workflow_name
1204 ))
1205 })?
1206 .clone();
1207
1208 #[cfg(feature = "prometheus")]
1209 gauge!(RUNS_ACTIVE, "workflow" => run.workflow_name.clone()).increment(1.0);
1210
1211 let run_start = Instant::now();
1212 let mut ctx = self.build_context_with_guard(&run, handler.as_ref());
1213
1214 let result = if handler.is_version_compatible(run.handler_version.as_deref()) {
1227 ctx.load_replay_steps().await?;
1228 self.release_then_execute(run_id, handler.as_ref(), &mut ctx)
1229 .await
1230 } else {
1231 Err(EngineError::HandlerVersionMismatch {
1232 run_id,
1233 workflow_name: run.workflow_name.clone(),
1234 run_version: run
1235 .handler_version
1236 .clone()
1237 .unwrap_or_else(|| "unknown".to_string()),
1238 current_version: handler
1239 .version()
1240 .map(str::to_string)
1241 .unwrap_or_else(|| "unknown".to_string()),
1242 })
1243 };
1244
1245 self.finalize_run(
1246 run_id,
1247 &run.workflow_name,
1248 result,
1249 &ctx,
1250 run_start,
1251 run.labels,
1252 )
1253 .await
1254 }
1255
1256 #[tracing::instrument(name = "engine.execute_run", skip_all, fields(run_id = %run_id))]
1264 pub async fn execute_run(&self, run_id: Uuid) -> Result<WorkflowResult, EngineError> {
1265 self.execute_handler_run(run_id).await
1266 }
1267
1268 #[tracing::instrument(name = "engine.resume_run", skip_all, fields(run_id = %run_id))]
1297 pub async fn resume_run(&self, run_id: Uuid) -> Result<WorkflowResult, EngineError> {
1298 let run = self
1299 .store
1300 .get_run(run_id)
1301 .await?
1302 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
1303
1304 if let Some(root_run_id) = chain_root(&run) {
1305 return self.resume_chain(run, root_run_id).await;
1306 }
1307
1308 self.resume_loaded_run(run).await
1309 }
1310
1311 async fn resume_chain(
1326 &self,
1327 child: Run,
1328 root_run_id: Uuid,
1329 ) -> Result<WorkflowResult, EngineError> {
1330 let child_run_id = child.id;
1331 let lease = child.worker_id.zip(child.lease_expires_at);
1332 let root = self
1333 .store
1334 .get_run(root_run_id)
1335 .await?
1336 .ok_or(EngineError::Store(StoreError::RunNotFound(root_run_id)))?;
1337
1338 match root.status.state {
1339 RunStatus::AwaitingApproval | RunStatus::Pending => {
1340 self.move_root_to_running(root_run_id, lease.as_ref())
1341 .await?;
1342 }
1343 RunStatus::Sleeping => {
1344 self.store
1345 .update_run_status(root_run_id, RunStatus::Pending)
1346 .await?;
1347 self.move_root_to_running(root_run_id, lease.as_ref())
1348 .await?;
1349 }
1350 other => {
1351 let reason = format!(
1352 "cannot resume child run {child_run_id}: root run {root_run_id} is {other}"
1353 );
1354 if let Err(err) = self
1355 .fail_or_schedule_retry(child_run_id, &reason, false, None, None)
1356 .await
1357 {
1358 error!(
1359 run_id = %child_run_id,
1360 error = %err,
1361 "failed to fail a child run whose root cannot resume"
1362 );
1363 }
1364 return Err(EngineError::InvalidWorkflow(reason));
1365 }
1366 }
1367
1368 if lease.is_some() {
1369 self.store
1370 .update_run(
1371 child_run_id,
1372 RunUpdate {
1373 lease: Some(LeaseUpdate::Release),
1374 ..RunUpdate::default()
1375 },
1376 )
1377 .await?;
1378 }
1379
1380 info!(
1381 run_id = %child_run_id,
1382 root_run_id = %root_run_id,
1383 lease_transferred = lease.is_some(),
1384 "child run resumed through its root run"
1385 );
1386
1387 let root = self
1388 .store
1389 .get_run(root_run_id)
1390 .await?
1391 .ok_or(EngineError::Store(StoreError::RunNotFound(root_run_id)))?;
1392 self.resume_loaded_run(root).await
1393 }
1394
1395 async fn move_root_to_running(
1401 &self,
1402 root_run_id: Uuid,
1403 lease: Option<&(String, DateTime<Utc>)>,
1404 ) -> Result<(), EngineError> {
1405 match lease {
1406 Some((worker_id, expires_at)) => {
1407 self.store
1408 .update_run(
1409 root_run_id,
1410 RunUpdate {
1411 status: Some(RunStatus::Running),
1412 lease: Some(LeaseUpdate::Set {
1413 worker_id: worker_id.clone(),
1414 expires_at: *expires_at,
1415 }),
1416 ..RunUpdate::default()
1417 },
1418 )
1419 .await?;
1420 }
1421 None => {
1422 self.store
1423 .update_run_status(root_run_id, RunStatus::Running)
1424 .await?;
1425 }
1426 }
1427 Ok(())
1428 }
1429
1430 async fn resume_loaded_run(&self, run: Run) -> Result<WorkflowResult, EngineError> {
1432 let run_id = run.id;
1433 let handler = self
1434 .handlers
1435 .get(&run.workflow_name)
1436 .ok_or_else(|| {
1437 EngineError::InvalidWorkflow(format!(
1438 "no handler registered: {}",
1439 run.workflow_name
1440 ))
1441 })?
1442 .clone();
1443
1444 info!(run_id = %run_id, workflow = %run.workflow_name, "resuming run after approval");
1445
1446 let run_start = Instant::now();
1447 let mut ctx = self.build_context_with_guard(&run, handler.as_ref());
1448
1449 let result = if handler.is_version_compatible(run.handler_version.as_deref()) {
1450 ctx.load_replay_steps().await?;
1451 self.release_then_execute(run_id, handler.as_ref(), &mut ctx)
1452 .await
1453 } else {
1454 Err(EngineError::HandlerVersionMismatch {
1455 run_id,
1456 workflow_name: run.workflow_name.clone(),
1457 run_version: run
1458 .handler_version
1459 .clone()
1460 .unwrap_or_else(|| "unknown".to_string()),
1461 current_version: handler
1462 .version()
1463 .map(str::to_string)
1464 .unwrap_or_else(|| "unknown".to_string()),
1465 })
1466 };
1467
1468 self.finalize_run(
1469 run_id,
1470 &run.workflow_name,
1471 result,
1472 &ctx,
1473 run_start,
1474 run.labels,
1475 )
1476 .await
1477 }
1478
1479 pub async fn deliver_signal(
1525 self: &Arc<Self>,
1526 signal: NewSignal,
1527 ) -> Result<SignalDelivery, EngineError> {
1528 if signal.name.trim().is_empty() {
1529 return Err(EngineError::InvalidSignal(
1530 "signal name must not be empty".to_string(),
1531 ));
1532 }
1533 if signal.key.trim().is_empty() {
1534 return Err(EngineError::InvalidSignal(
1535 "signal key must not be empty".to_string(),
1536 ));
1537 }
1538
1539 let stored = match self.store.insert_signal(signal).await? {
1540 SignalInsert::Created(stored) => stored,
1541 SignalInsert::Duplicate(existing) => {
1542 info!(
1543 signal_id = %existing.id,
1544 signal = %existing.name,
1545 key = %existing.key,
1546 "duplicate signal ignored"
1547 );
1548 return Ok(SignalDelivery {
1549 signal_id: existing.id,
1550 duplicate: true,
1551 resumed: Vec::new(),
1552 rejected: Vec::new(),
1553 });
1554 }
1555 };
1556
1557 let waiters = self
1558 .store
1559 .list_signal_waiters(&stored.name, &stored.key)
1560 .await?;
1561 let mut resumed = Vec::new();
1562 let mut rejected = Vec::new();
1563
1564 for step in waiters {
1565 if let Err(error) = validate_step_payload(step.input.as_ref(), &stored.payload) {
1566 rejected.push(SignalRejected {
1567 run_id: step.run_id,
1568 step_id: step.id,
1569 error,
1570 });
1571 continue;
1572 }
1573
1574 match self
1575 .store
1576 .resolve_signal_step(step.id, received_output(&stored))
1577 .await
1578 {
1579 Ok(SignalStepResolution::Resolved {
1580 run_id,
1581 run_resumed,
1582 }) => {
1583 resumed.push(SignalResumed {
1584 run_id,
1585 step_id: step.id,
1586 });
1587 if run_resumed && self.execution_mode == ExecutionMode::Local {
1588 self.spawn_local_resume(run_id);
1589 }
1590 }
1591 Ok(SignalStepResolution::NotWaiting { .. }) => {}
1593 Err(err) => {
1594 error!(
1595 run_id = %step.run_id,
1596 step_id = %step.id,
1597 error = %err,
1598 "failed to resolve a waiting signal step"
1599 );
1600 rejected.push(SignalRejected {
1601 run_id: step.run_id,
1602 step_id: step.id,
1603 error: err.to_string(),
1604 });
1605 }
1606 }
1607 }
1608
1609 info!(
1610 signal_id = %stored.id,
1611 signal = %stored.name,
1612 key = %stored.key,
1613 resumed = resumed.len(),
1614 rejected = rejected.len(),
1615 "signal received"
1616 );
1617 self.event_publisher
1618 .publish(Event::SignalReceived(SignalReceivedEvent {
1619 signal_id: stored.id,
1620 name: stored.name.clone(),
1621 key: stored.key.clone(),
1622 resumed_runs: resumed.iter().map(|r| r.run_id).collect(),
1623 at: stored.received_at,
1624 }));
1625
1626 Ok(SignalDelivery {
1627 signal_id: stored.id,
1628 duplicate: false,
1629 resumed,
1630 rejected,
1631 })
1632 }
1633
1634 pub async fn send_signal<S: Signal>(
1672 self: &Arc<Self>,
1673 signal: &S,
1674 key: &str,
1675 idempotency_id: Option<&str>,
1676 ) -> Result<SignalDelivery, EngineError> {
1677 let payload = to_value(signal)?;
1678 self.deliver_signal(NewSignal {
1679 name: S::NAME.to_string(),
1680 key: key.to_string(),
1681 payload,
1682 idempotency_id: idempotency_id.map(str::to_string),
1683 })
1684 .await
1685 }
1686
1687 pub(crate) fn spawn_local_resume(self: &Arc<Self>, run_id: Uuid) {
1693 let engine = Arc::clone(self);
1694 spawn(async move {
1695 if let Err(err) = engine
1696 .store
1697 .update_run_status(run_id, RunStatus::Running)
1698 .await
1699 {
1700 error!(run_id = %run_id, error = %err, "failed to restart a woken run");
1701 return;
1702 }
1703 if let Err(err) = engine.resume_run(run_id).await {
1704 error!(run_id = %run_id, error = %err, "failed to resume a woken run");
1705 }
1706 });
1707 }
1708
1709 pub async fn fail_or_schedule_retry(
1757 &self,
1758 run_id: Uuid,
1759 error: &str,
1760 retryable: bool,
1761 cost_usd: Option<Decimal>,
1762 duration_ms: Option<u64>,
1763 ) -> Result<RunStatus, EngineError> {
1764 let run = self
1765 .store
1766 .get_run(run_id)
1767 .await?
1768 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
1769
1770 let has_attempts_left = run.retry_count < run.max_retries;
1771 let update = if retryable && has_attempts_left {
1772 let backoff = backoff_for_retry(run.retry_count);
1773 let scheduled_at = Utc::now() + TimeDelta::milliseconds(backoff.as_millis() as i64);
1774
1775 info!(
1776 run_id = %run_id,
1777 workflow = %run.workflow_name,
1778 attempt = run.retry_count + 1,
1779 max_retries = run.max_retries,
1780 backoff_secs = backoff.as_secs(),
1781 scheduled_at = %scheduled_at,
1782 "run failed, scheduling retry"
1783 );
1784
1785 RunUpdate {
1786 status: Some(RunStatus::Retrying),
1787 error: Some(error.to_string()),
1788 increment_retry: true,
1789 cost_usd,
1790 duration_ms,
1791 scheduled_at: Some(scheduled_at),
1792 ..RunUpdate::default()
1793 }
1794 } else {
1795 RunUpdate {
1796 status: Some(RunStatus::Failed),
1797 error: Some(error.to_string()),
1798 cost_usd,
1799 duration_ms,
1800 completed_at: Some(Utc::now()),
1801 ..RunUpdate::default()
1802 }
1803 };
1804
1805 let status = update.status.unwrap_or(RunStatus::Failed);
1806 self.store.update_run(run_id, update).await?;
1807 self.fail_orphaned_steps(run_id, error).await?;
1808 self.cancel_descendants_of_stopped_run(run_id, error).await;
1811
1812 Ok(status)
1813 }
1814
1815 pub async fn interrupt_running_steps(&self, run_id: Uuid) -> Result<(), EngineError> {
1845 interrupt_running_steps(self.store.as_ref(), run_id).await
1846 }
1847
1848 pub async fn fail_orphaned_steps(
1862 &self,
1863 run_id: Uuid,
1864 error_message: &str,
1865 ) -> Result<(), EngineError> {
1866 let steps = self.store.list_steps(run_id).await?;
1867 let now = Utc::now();
1868
1869 for step in steps {
1870 if step.status.state.is_terminal() {
1871 continue;
1872 }
1873
1874 let (target_status, error) = match step.status.state {
1875 StepStatus::Running | StepStatus::AwaitingApproval => {
1876 let err = if step.error.is_some() {
1877 None
1878 } else {
1879 Some(error_message.to_string())
1880 };
1881 (StepStatus::Failed, err)
1882 }
1883 StepStatus::Pending => (StepStatus::Skipped, None),
1884 _ => continue,
1885 };
1886
1887 if let Err(e) = self
1888 .store
1889 .update_step(
1890 step.id,
1891 StepUpdate {
1892 status: Some(target_status),
1893 error,
1894 completed_at: Some(now),
1895 ..StepUpdate::default()
1896 },
1897 )
1898 .await
1899 {
1900 warn!(
1901 run_id = %run_id,
1902 step_id = %step.id,
1903 step_name = %step.name,
1904 error = %e,
1905 "failed to cleanup orphaned step"
1906 );
1907 } else {
1908 info!(
1909 run_id = %run_id,
1910 step_id = %step.id,
1911 step_name = %step.name,
1912 from = %step.status.state,
1913 to = %target_status,
1914 "cleaned up orphaned step"
1915 );
1916 }
1917 }
1918
1919 Ok(())
1920 }
1921
1922 async fn release_then_execute(
1927 &self,
1928 run_id: Uuid,
1929 handler: &dyn WorkflowHandler,
1930 ctx: &mut WorkflowContext,
1931 ) -> Result<(), EngineError> {
1932 match self.provider.release_run(&run_id.to_string()).await {
1933 Ok(()) => handler.execute(ctx).await,
1934 Err(e) => Err(EngineError::Operation(OperationError::Agent(e))),
1935 }
1936 }
1937
1938 async fn finalize_run(
1944 &self,
1945 run_id: Uuid,
1946 workflow_name: &str,
1947 result: Result<(), EngineError>,
1948 ctx: &WorkflowContext,
1949 run_start: Instant,
1950 run_labels: HashMap<String, String>,
1951 ) -> Result<WorkflowResult, EngineError> {
1952 let total_duration = ctx.carried_duration_ms() + run_start.elapsed().as_millis() as u64;
1955 let completed_at = Utc::now();
1956
1957 let final_status;
1958 let final_run;
1959
1960 match result {
1961 Ok(()) => {
1962 final_status = if ctx.has_allowed_failure() {
1963 RunStatus::Warning
1964 } else {
1965 RunStatus::Completed
1966 };
1967 final_run = self
1968 .store
1969 .update_run_returning(
1970 run_id,
1971 RunUpdate {
1972 status: Some(final_status),
1973 cost_usd: Some(ctx.total_cost_usd()),
1974 duration_ms: Some(total_duration),
1975 completed_at: Some(completed_at),
1976 output: ctx.output().cloned(),
1977 ..RunUpdate::default()
1978 },
1979 )
1980 .await?;
1981
1982 info!(
1983 run_id = %run_id,
1984 status = %final_status,
1985 cost_usd = %ctx.total_cost_usd(),
1986 duration_ms = total_duration,
1987 "run completed"
1988 );
1989 }
1990 Err(EngineError::ApprovalRequired {
1991 run_id: approval_run_id,
1992 step_id,
1993 ref message,
1994 }) => {
1995 final_status = RunStatus::AwaitingApproval;
1996 final_run = self
1997 .store
1998 .update_run_returning(
1999 run_id,
2000 RunUpdate {
2001 status: Some(RunStatus::AwaitingApproval),
2002 cost_usd: Some(ctx.total_cost_usd()),
2003 duration_ms: Some(total_duration),
2004 ..RunUpdate::default()
2005 },
2006 )
2007 .await?;
2008
2009 info!(
2010 run_id = %approval_run_id,
2011 step_id = %step_id,
2012 message = %message,
2013 "run awaiting approval"
2014 );
2015
2016 self.publish_approval_requested(approval_run_id, step_id, message)
2017 .await?;
2018 }
2019 Err(EngineError::ChildSuspended {
2020 run_id: child_run_id,
2021 ref cause,
2022 }) => {
2023 final_status = cause.suspension_status();
2024 final_run = self
2028 .store
2029 .update_run_returning(
2030 run_id,
2031 RunUpdate {
2032 status: Some(final_status),
2033 cost_usd: Some(ctx.total_cost_usd()),
2034 duration_ms: Some(total_duration),
2035 ..RunUpdate::default()
2036 },
2037 )
2038 .await?;
2039
2040 let leaf = cause.suspension_leaf();
2041 info!(
2042 run_id = %run_id,
2043 child_run_id = %child_run_id,
2044 status = %final_status,
2045 cause = %leaf,
2046 "run suspended with its child run"
2047 );
2048
2049 match leaf {
2050 EngineError::ApprovalRequired {
2051 run_id: approval_run_id,
2052 step_id,
2053 message,
2054 } => {
2055 self.publish_approval_requested(*approval_run_id, *step_id, message)
2056 .await?;
2057 }
2058 EngineError::SignalWaiting {
2059 run_id: wait_run_id,
2060 step_id,
2061 step_name,
2062 name,
2063 key,
2064 deadline_at,
2065 } => {
2066 self.event_publisher
2067 .publish(Event::SignalAwaited(SignalAwaitedEvent {
2068 run_id: *wait_run_id,
2069 step_id: *step_id,
2070 step_name: step_name.clone(),
2071 name: name.clone(),
2072 key: key.clone(),
2073 deadline_at: *deadline_at,
2074 at: Utc::now(),
2075 }));
2076 }
2077 _ => {}
2080 }
2081 }
2082 Err(EngineError::HumanInputRequired {
2083 run_id: input_run_id,
2084 step_id,
2085 ref message,
2086 }) => {
2087 final_status = RunStatus::AwaitingApproval;
2088 final_run = self
2089 .store
2090 .update_run_returning(
2091 run_id,
2092 RunUpdate {
2093 status: Some(RunStatus::AwaitingApproval),
2094 cost_usd: Some(ctx.total_cost_usd()),
2095 duration_ms: Some(total_duration),
2096 ..RunUpdate::default()
2097 },
2098 )
2099 .await?;
2100
2101 info!(
2103 run_id = %input_run_id,
2104 step_id = %step_id,
2105 message = %message,
2106 "run awaiting human input"
2107 );
2108 }
2109 Err(EngineError::DelaySleeping {
2110 run_id: delay_run_id,
2111 step_id,
2112 wake_at,
2113 }) => {
2114 final_status = RunStatus::Sleeping;
2115 final_run = self
2116 .store
2117 .update_run_returning(
2118 run_id,
2119 RunUpdate {
2120 status: Some(RunStatus::Sleeping),
2121 cost_usd: Some(ctx.total_cost_usd()),
2122 duration_ms: Some(total_duration),
2123 scheduled_at: Some(wake_at),
2124 ..RunUpdate::default()
2125 },
2126 )
2127 .await?;
2128
2129 info!(
2130 run_id = %delay_run_id,
2131 step_id = %step_id,
2132 wake_at = %wake_at,
2133 "run sleeping until delay elapses"
2134 );
2135 }
2136 Err(EngineError::CapacitySleeping {
2137 run_id: capacity_run_id,
2138 step_id,
2139 ref kind,
2140 wake_at,
2141 }) => {
2142 final_status = RunStatus::Sleeping;
2143 final_run = self
2144 .store
2145 .update_run_returning(
2146 run_id,
2147 RunUpdate {
2148 status: Some(RunStatus::Sleeping),
2149 cost_usd: Some(ctx.total_cost_usd()),
2150 duration_ms: Some(total_duration),
2151 scheduled_at: Some(wake_at),
2152 capacity_wait_kind: Some(ProviderKind::new(kind.as_str())),
2153 ..RunUpdate::default()
2154 },
2155 )
2156 .await?;
2157
2158 info!(
2159 run_id = %capacity_run_id,
2160 step_id = %step_id,
2161 kind = %kind,
2162 wake_at = %wake_at,
2163 "run sleeping until provider capacity returns"
2164 );
2165 }
2166 Err(EngineError::SignalWaiting {
2167 run_id: wait_run_id,
2168 step_id,
2169 ref step_name,
2170 ref name,
2171 ref key,
2172 deadline_at,
2173 }) => {
2174 final_status = RunStatus::Sleeping;
2175 let waiting = self
2179 .store
2180 .suspend_run_on_signal(run_id, step_id, deadline_at)
2181 .await?;
2182 final_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
2194 if waiting {
2195 self.event_publisher
2196 .publish(Event::SignalAwaited(SignalAwaitedEvent {
2197 run_id: wait_run_id,
2198 step_id,
2199 step_name: step_name.clone(),
2200 name: name.clone(),
2201 key: key.clone(),
2202 deadline_at,
2203 at: Utc::now(),
2204 }));
2205 }
2206
2207 info!(
2208 run_id = %wait_run_id,
2209 step_id = %step_id,
2210 signal = %name,
2211 key = %key,
2212 deadline_at = %deadline_at,
2213 waiting,
2214 "run sleeping until a signal arrives"
2215 );
2216 }
2217 Err(err) => {
2218 let guardrail_stop = matches!(
2222 err,
2223 EngineError::RunBudgetExceeded { .. } | EngineError::WorkflowGuardRejected(_)
2224 );
2225
2226 final_status = if guardrail_stop {
2227 if let Err(store_err) = self
2228 .store
2229 .update_run(
2230 run_id,
2231 RunUpdate {
2232 status: Some(RunStatus::Cancelled),
2233 error: Some(err.to_string()),
2234 cost_usd: Some(ctx.total_cost_usd()),
2235 duration_ms: Some(total_duration),
2236 completed_at: Some(completed_at),
2237 output: ctx.output().cloned(),
2238 ..RunUpdate::default()
2239 },
2240 )
2241 .await
2242 {
2243 error!(run_id = %run_id, store_error = %store_err, "failed to persist run cancellation");
2244 }
2245 if let Err(cleanup_err) = self
2246 .fail_orphaned_steps(run_id, "run stopped: guardrail limit reached")
2247 .await
2248 {
2249 error!(run_id = %run_id, store_error = %cleanup_err, "failed to cleanup orphaned steps");
2250 }
2251 RunStatus::Cancelled
2252 } else {
2253 if let Some(output) = ctx.output()
2256 && let Err(store_err) = self
2257 .store
2258 .update_run(
2259 run_id,
2260 RunUpdate {
2261 output: Some(output.clone()),
2262 ..RunUpdate::default()
2263 },
2264 )
2265 .await
2266 {
2267 error!(run_id = %run_id, store_error = %store_err, "failed to persist run output");
2268 }
2269 self.fail_or_schedule_retry(
2270 run_id,
2271 &err.to_string(),
2272 is_run_retryable(&err),
2273 Some(ctx.total_cost_usd()),
2274 Some(total_duration),
2275 )
2276 .await
2277 .unwrap_or_else(|store_err| {
2278 error!(run_id = %run_id, store_error = %store_err, "failed to persist run failure");
2279 RunStatus::Failed
2280 })
2281 };
2282
2283 if matches!(err, EngineError::RunBudgetExceeded { .. }) {
2284 self.on_run_budget_exceeded(workflow_name, run_id, &err);
2285 }
2286
2287 error!(run_id = %run_id, status = %final_status, error = %err, "run stopped");
2288
2289 self.publish_run_status_changed(
2290 workflow_name,
2291 run_id,
2292 final_status,
2293 Some(err.to_string()),
2294 ctx,
2295 total_duration,
2296 run_labels,
2297 );
2298
2299 #[cfg(feature = "prometheus")]
2300 self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
2301
2302 return Err(err);
2303 }
2304 }
2305
2306 self.publish_run_status_changed(
2307 workflow_name,
2308 run_id,
2309 final_status,
2310 None,
2311 ctx,
2312 total_duration,
2313 run_labels,
2314 );
2315
2316 #[cfg(feature = "prometheus")]
2317 self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
2318
2319 Ok(WorkflowResult {
2320 run: final_run,
2321 steps: ctx.step_results().to_vec(),
2322 })
2323 }
2324
2325 async fn publish_approval_requested(
2328 &self,
2329 run_id: Uuid,
2330 step_id: Uuid,
2331 message: &str,
2332 ) -> Result<(), EngineError> {
2333 let requirement = self
2334 .store
2335 .get_step(step_id)
2336 .await?
2337 .and_then(|s| s.approval_requirement);
2338 self.event_publisher
2339 .publish(Event::ApprovalRequested(ApprovalRequestedEvent {
2340 run_id,
2341 step_id,
2342 message: message.to_string(),
2343 requirement,
2344 at: Utc::now(),
2345 }));
2346 Ok(())
2347 }
2348
2349 pub async fn fail_ancestors(&self, run_id: Uuid, reason: &str) -> Result<(), EngineError> {
2378 let mut current = self
2379 .store
2380 .get_run(run_id)
2381 .await?
2382 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
2383 let mut visited = HashSet::from([run_id]);
2385
2386 while let Some(parent_id) = chain_parent(¤t) {
2387 if !visited.insert(parent_id) {
2388 break;
2389 }
2390 let status = self
2391 .fail_or_schedule_retry(parent_id, reason, false, None, None)
2392 .await?;
2393 info!(
2394 run_id = %run_id,
2395 ancestor_run_id = %parent_id,
2396 status = %status,
2397 "ancestor run failed with its child"
2398 );
2399 current = self
2400 .store
2401 .get_run(parent_id)
2402 .await?
2403 .ok_or(EngineError::Store(StoreError::RunNotFound(parent_id)))?;
2404 }
2405
2406 Ok(())
2407 }
2408
2409 #[cfg(feature = "prometheus")]
2411 fn emit_run_metrics(
2412 &self,
2413 workflow_name: &str,
2414 status: RunStatus,
2415 duration_ms: u64,
2416 ctx: &WorkflowContext,
2417 ) {
2418 let status_str = status.to_string();
2419 let wf = workflow_name.to_string();
2420
2421 counter!(RUNS_TOTAL, "workflow" => wf.clone(), "status" => status_str.clone()).increment(1);
2422 histogram!(RUN_DURATION_SECONDS, "workflow" => wf.clone(), "status" => status_str)
2423 .record(duration_ms as f64 / 1000.0);
2424 histogram!(RUN_COST_USD, "workflow" => wf.clone()).record(
2425 ctx.total_cost_usd()
2426 .to_string()
2427 .parse::<f64>()
2428 .unwrap_or(0.0),
2429 );
2430 gauge!(RUNS_ACTIVE, "workflow" => wf).decrement(1.0);
2431 }
2432
2433 fn on_run_budget_exceeded(&self, workflow_name: &str, run_id: Uuid, err: &EngineError) {
2439 let EngineError::RunBudgetExceeded {
2440 limit_usd,
2441 spent_usd,
2442 step_budget_usd,
2443 ..
2444 } = err
2445 else {
2446 return;
2447 };
2448
2449 #[cfg(feature = "prometheus")]
2450 counter!(
2451 RUN_BUDGET_EXCEEDED_TOTAL,
2452 "workflow" => workflow_name.to_string(),
2453 "scope" => "run",
2454 )
2455 .increment(1);
2456
2457 self.event_publisher
2458 .publish(Event::RunBudgetExceeded(RunBudgetExceededEvent {
2459 run_id,
2460 workflow_name: workflow_name.to_string(),
2461 limit_usd: *limit_usd,
2462 spent_usd: *spent_usd,
2463 step_budget_usd: *step_budget_usd,
2464 at: Utc::now(),
2465 }));
2466 }
2467
2468 #[allow(clippy::too_many_arguments)]
2473 fn publish_run_status_changed(
2474 &self,
2475 workflow_name: &str,
2476 run_id: Uuid,
2477 to: RunStatus,
2478 error: Option<String>,
2479 ctx: &WorkflowContext,
2480 duration_ms: u64,
2481 labels: HashMap<String, String>,
2482 ) {
2483 let now = Utc::now();
2484 let cost_usd = ctx.total_cost_usd();
2485 let wf = workflow_name.to_string();
2486
2487 self.event_publisher
2488 .publish(Event::RunStatusChanged(RunStatusChangedEvent {
2489 run_id,
2490 workflow_name: wf.clone(),
2491 from: RunStatus::Running,
2492 to,
2493 error: error.clone(),
2494 cost_usd,
2495 duration_ms,
2496 labels: labels.clone(),
2497 at: now,
2498 }));
2499
2500 if to == RunStatus::Failed {
2501 self.event_publisher
2502 .publish(Event::RunFailed(RunFailedEvent {
2503 run_id,
2504 workflow_name: wf,
2505 error,
2506 cost_usd,
2507 duration_ms,
2508 labels,
2509 at: now,
2510 }));
2511 }
2512 }
2513}
2514
2515impl fmt::Debug for Engine {
2516 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
2517 f.debug_struct("Engine")
2518 .field("handlers", &self.handlers.keys().collect::<Vec<_>>())
2519 .finish_non_exhaustive()
2520 }
2521}
2522
2523#[cfg(test)]
2524mod tests {
2525 use super::*;
2526 use crate::config::ShellConfig;
2527 use crate::handler::{HandlerFuture, WorkflowHandler};
2528 use ironflow_core::providers::claude::ClaudeCodeProvider;
2529 use ironflow_core::providers::record_replay::RecordReplayProvider;
2530 use ironflow_store::memory::InMemoryStore;
2531 use ironflow_store::models::StepStatus;
2532 use serde_json::json;
2533
2534 struct EchoWorkflow;
2536
2537 impl WorkflowHandler for EchoWorkflow {
2538 fn name(&self) -> &str {
2539 "echo-workflow"
2540 }
2541
2542 fn describe(&self) -> WorkflowInfo {
2543 WorkflowInfo {
2544 description: "A simple workflow that echoes hello".to_string(),
2545 source_code: None,
2546 sub_workflows: Vec::new(),
2547 category: None,
2548 version: self.version().map(str::to_string),
2549 compatible_versions: Vec::new(),
2550 input_schema: None,
2551 default_labels: HashMap::new(),
2552 schedule: self.schedule().cloned(),
2553 default_max_cost_usd: self.default_max_cost_usd(),
2554 }
2555 }
2556
2557 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2558 Box::pin(async move {
2559 ctx.shell("greet", ShellConfig::new("echo hello")).await?;
2560 Ok(())
2561 })
2562 }
2563 }
2564
2565 struct FailingWorkflow;
2567
2568 impl WorkflowHandler for FailingWorkflow {
2569 fn name(&self) -> &str {
2570 "failing-workflow"
2571 }
2572
2573 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2574 Box::pin(async move {
2575 ctx.shell("fail", ShellConfig::new("exit 1")).await?;
2576 Ok(())
2577 })
2578 }
2579 }
2580
2581 fn create_test_engine() -> Engine {
2582 let store = Arc::new(InMemoryStore::new());
2583 let inner = ClaudeCodeProvider::new();
2584 let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
2585 inner,
2586 "/tmp/ironflow-fixtures",
2587 ));
2588 Engine::new(store, provider)
2589 }
2590
2591 #[test]
2592 fn engine_new_creates_instance() {
2593 let engine = create_test_engine();
2594 assert_eq!(engine.handler_names().len(), 0);
2595 }
2596
2597 #[test]
2598 fn execution_mode_defaults_to_local() {
2599 let engine = create_test_engine();
2600 assert_eq!(engine.execution_mode(), ExecutionMode::Local);
2601 }
2602
2603 #[test]
2604 fn with_execution_mode_overrides_the_default() {
2605 let engine = create_test_engine().with_execution_mode(ExecutionMode::Workers);
2606 assert_eq!(engine.execution_mode(), ExecutionMode::Workers);
2607 }
2608
2609 #[test]
2610 fn engine_register_handler() {
2611 let mut engine = create_test_engine();
2612 let result = engine.register(EchoWorkflow);
2613 assert!(result.is_ok());
2614 assert_eq!(engine.handler_names().len(), 1);
2615 assert!(engine.handler_names().contains(&"echo-workflow"));
2616 }
2617
2618 #[test]
2619 fn engine_register_duplicate_returns_error() {
2620 let mut engine = create_test_engine();
2621 engine.register(EchoWorkflow).unwrap();
2622 let result = engine.register(EchoWorkflow);
2623 assert!(result.is_err());
2624 }
2625
2626 #[test]
2627 fn engine_get_handler_found() {
2628 let mut engine = create_test_engine();
2629 engine.register(EchoWorkflow).unwrap();
2630 let handler = engine.get_handler("echo-workflow");
2631 assert!(handler.is_some());
2632 }
2633
2634 #[test]
2635 fn engine_get_handler_not_found() {
2636 let engine = create_test_engine();
2637 let handler = engine.get_handler("nonexistent");
2638 assert!(handler.is_none());
2639 }
2640
2641 #[test]
2642 fn engine_handler_names_lists_all() {
2643 let mut engine = create_test_engine();
2644 engine.register(EchoWorkflow).unwrap();
2645 engine.register(FailingWorkflow).unwrap();
2646 let names = engine.handler_names();
2647 assert_eq!(names.len(), 2);
2648 assert!(names.contains(&"echo-workflow"));
2649 assert!(names.contains(&"failing-workflow"));
2650 }
2651
2652 #[test]
2653 fn engine_handler_info_returns_description() {
2654 let mut engine = create_test_engine();
2655 engine.register(EchoWorkflow).unwrap();
2656 let info = engine.handler_info("echo-workflow");
2657 assert!(info.is_some());
2658 let info = info.unwrap();
2659 assert_eq!(info.description, "A simple workflow that echoes hello");
2660 }
2661
2662 struct CategorizedWorkflow;
2663
2664 impl WorkflowHandler for CategorizedWorkflow {
2665 fn name(&self) -> &str {
2666 "categorized"
2667 }
2668 fn category(&self) -> Option<&str> {
2669 Some("data/etl")
2670 }
2671 fn execute<'a>(
2672 &'a self,
2673 _ctx: &'a mut WorkflowContext,
2674 ) -> crate::handler::HandlerFuture<'a> {
2675 Box::pin(async move { Ok(()) })
2676 }
2677 }
2678
2679 #[test]
2680 fn engine_default_describe_propagates_category() {
2681 let mut engine = create_test_engine();
2682 engine.register(CategorizedWorkflow).unwrap();
2683 let info = engine.handler_info("categorized").unwrap();
2684 assert_eq!(info.category.as_deref(), Some("data/etl"));
2685 }
2686
2687 #[test]
2688 fn engine_default_describe_without_category() {
2689 let mut engine = create_test_engine();
2690 engine.register(EchoWorkflow).unwrap();
2691 let info = engine.handler_info("echo-workflow").unwrap();
2692 assert!(info.category.is_none());
2693 }
2694
2695 struct ScheduledWorkflow {
2700 schedule: CronSchedule,
2701 }
2702
2703 impl ScheduledWorkflow {
2704 fn new() -> Self {
2705 Self {
2706 schedule: CronSchedule::new("0 0 * * * *").unwrap(),
2707 }
2708 }
2709 }
2710
2711 impl WorkflowHandler for ScheduledWorkflow {
2712 fn name(&self) -> &str {
2713 "scheduled"
2714 }
2715 fn schedule(&self) -> Option<&CronSchedule> {
2716 Some(&self.schedule)
2717 }
2718 fn execute<'a>(
2719 &'a self,
2720 _ctx: &'a mut WorkflowContext,
2721 ) -> crate::handler::HandlerFuture<'a> {
2722 Box::pin(async move { Ok(()) })
2723 }
2724 }
2725
2726 #[test]
2727 fn engine_default_describe_propagates_schedule() {
2728 let mut engine = create_test_engine();
2729 engine.register(ScheduledWorkflow::new()).unwrap();
2730 let info = engine.handler_info("scheduled").unwrap();
2731 assert_eq!(
2732 info.schedule.as_ref().map(|s| s.as_str()),
2733 Some("0 0 * * * *")
2734 );
2735 }
2736
2737 #[test]
2738 fn engine_default_describe_without_schedule() {
2739 let mut engine = create_test_engine();
2740 engine.register(EchoWorkflow).unwrap();
2741 let info = engine.handler_info("echo-workflow").unwrap();
2742 assert!(info.schedule.is_none());
2743 }
2744
2745 #[test]
2746 fn scheduled_handlers_returns_only_scheduled() {
2747 let mut engine = create_test_engine();
2748 engine.register(EchoWorkflow).unwrap();
2749 engine.register(ScheduledWorkflow::new()).unwrap();
2750 engine.register(FailingWorkflow).unwrap();
2751
2752 let scheduled = engine.scheduled_handlers();
2753 assert_eq!(scheduled.len(), 1);
2754 assert_eq!(scheduled[0].0, "scheduled");
2755 assert_eq!(scheduled[0].1.as_str(), "0 0 * * * *");
2756 }
2757
2758 #[test]
2759 fn scheduled_handlers_empty_when_none_scheduled() {
2760 let mut engine = create_test_engine();
2761 engine.register(EchoWorkflow).unwrap();
2762 engine.register(FailingWorkflow).unwrap();
2763
2764 let scheduled = engine.scheduled_handlers();
2765 assert!(scheduled.is_empty());
2766 }
2767
2768 struct BadCategoryWorkflow(&'static str);
2769
2770 impl WorkflowHandler for BadCategoryWorkflow {
2771 fn name(&self) -> &str {
2772 "bad-category"
2773 }
2774 fn category(&self) -> Option<&str> {
2775 Some(self.0)
2776 }
2777 fn execute<'a>(
2778 &'a self,
2779 _ctx: &'a mut WorkflowContext,
2780 ) -> crate::handler::HandlerFuture<'a> {
2781 Box::pin(async move { Ok(()) })
2782 }
2783 }
2784
2785 #[test]
2786 fn engine_register_rejects_empty_category() {
2787 let mut engine = create_test_engine();
2788 let err = engine.register(BadCategoryWorkflow("")).unwrap_err();
2789 match err {
2790 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty category")),
2791 other => panic!("expected InvalidWorkflow, got {other:?}"),
2792 }
2793 }
2794
2795 #[test]
2796 fn engine_register_rejects_leading_slash_category() {
2797 let mut engine = create_test_engine();
2798 let err = engine
2799 .register(BadCategoryWorkflow("/data/etl"))
2800 .unwrap_err();
2801 match err {
2802 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("leading '/'")),
2803 other => panic!("expected InvalidWorkflow, got {other:?}"),
2804 }
2805 }
2806
2807 #[test]
2808 fn engine_register_rejects_trailing_slash_category() {
2809 let mut engine = create_test_engine();
2810 let err = engine
2811 .register(BadCategoryWorkflow("data/etl/"))
2812 .unwrap_err();
2813 match err {
2814 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("trailing '/'")),
2815 other => panic!("expected InvalidWorkflow, got {other:?}"),
2816 }
2817 }
2818
2819 #[test]
2820 fn engine_register_rejects_double_slash_category() {
2821 let mut engine = create_test_engine();
2822 let err = engine
2823 .register(BadCategoryWorkflow("data//etl"))
2824 .unwrap_err();
2825 match err {
2826 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty segment")),
2827 other => panic!("expected InvalidWorkflow, got {other:?}"),
2828 }
2829 }
2830
2831 #[test]
2832 fn engine_register_rejects_whitespace_only_segment_category() {
2833 let mut engine = create_test_engine();
2834 let err = engine
2835 .register(BadCategoryWorkflow("data/ /etl"))
2836 .unwrap_err();
2837 match err {
2838 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("whitespace-only segment")),
2839 other => panic!("expected InvalidWorkflow, got {other:?}"),
2840 }
2841 }
2842
2843 #[test]
2844 fn engine_register_accepts_valid_nested_category() {
2845 let mut engine = create_test_engine();
2846 assert!(engine.register(CategorizedWorkflow).is_ok());
2847 }
2848
2849 #[tokio::test]
2850 async fn engine_unknown_workflow_returns_error() {
2851 let engine = create_test_engine();
2852 let result = engine
2853 .run_handler("unknown", TriggerKind::Manual, json!({}))
2854 .await;
2855 assert!(result.is_err());
2856 match result {
2857 Err(EngineError::InvalidWorkflow(msg)) => {
2858 assert!(msg.contains("no handler registered"));
2859 }
2860 _ => panic!("expected InvalidWorkflow error"),
2861 }
2862 }
2863
2864 #[tokio::test]
2865 async fn engine_enqueue_handler_creates_pending_run() {
2866 let mut engine = create_test_engine();
2867 engine.register(EchoWorkflow).unwrap();
2868
2869 let run = engine
2870 .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
2871 .await
2872 .unwrap();
2873 assert_eq!(run.status.state, RunStatus::Pending);
2874 assert_eq!(run.workflow_name, "echo-workflow");
2875 }
2876
2877 #[tokio::test]
2878 async fn enqueue_handler_leaves_the_run_unattributed() {
2879 let mut engine = create_test_engine();
2880 engine.register(EchoWorkflow).unwrap();
2881
2882 let run = engine
2883 .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
2884 .await
2885 .unwrap();
2886
2887 assert!(run.created_by.is_none());
2888 }
2889
2890 #[tokio::test]
2891 async fn enqueue_handler_with_options_records_the_author() {
2892 let mut engine = create_test_engine();
2893 engine.register(EchoWorkflow).unwrap();
2894 let actor = RunActor::User {
2895 user_id: Uuid::now_v7(),
2896 };
2897
2898 let run = engine
2899 .enqueue_handler_with_options(
2900 "echo-workflow",
2901 TriggerKind::Api,
2902 json!({}),
2903 EnqueueOptions {
2904 created_by: Some(actor.clone()),
2905 ..Default::default()
2906 },
2907 )
2908 .await
2909 .unwrap()
2910 .into_run();
2911
2912 assert_eq!(run.created_by, Some(actor));
2913 }
2914
2915 #[tokio::test]
2916 async fn enqueue_handler_with_options_accepts_no_author() {
2917 let mut engine = create_test_engine();
2918 engine.register(EchoWorkflow).unwrap();
2919
2920 let run = engine
2921 .enqueue_handler_with_options(
2922 "echo-workflow",
2923 TriggerKind::Cron {
2924 schedule: "0 * * * * *".to_string(),
2925 },
2926 json!({}),
2927 EnqueueOptions::default(),
2928 )
2929 .await
2930 .unwrap()
2931 .into_run();
2932
2933 assert!(run.created_by.is_none());
2934 }
2935
2936 #[tokio::test]
2937 async fn enqueue_handler_with_options_stores_concurrency_limits() {
2938 let mut engine = create_test_engine();
2939 engine.register(EchoWorkflow).unwrap();
2940 let limits = vec![
2941 ConcurrencyLimit::new("repo:acme", 2),
2942 ConcurrencyLimit::new("tenant:42", 5),
2943 ];
2944
2945 let run = engine
2946 .enqueue_handler_with_options(
2947 "echo-workflow",
2948 TriggerKind::Api,
2949 json!({}),
2950 EnqueueOptions {
2951 concurrency_limits: limits.clone(),
2952 ..Default::default()
2953 },
2954 )
2955 .await
2956 .unwrap()
2957 .into_run();
2958
2959 assert_eq!(run.concurrency_limits, limits);
2960 }
2961
2962 #[tokio::test]
2963 async fn enqueue_rejects_invalid_concurrency_limits() {
2964 let mut engine = create_test_engine();
2965 engine.register(EchoWorkflow).unwrap();
2966
2967 let invalid = [
2968 vec![ConcurrencyLimit::new("repo:acme", 0)],
2969 vec![ConcurrencyLimit::new("", 1)],
2970 vec![
2971 ConcurrencyLimit::new("repo:acme", 1),
2972 ConcurrencyLimit::new("repo:acme", 2),
2973 ],
2974 ];
2975 for concurrency_limits in invalid {
2976 let err = engine
2977 .enqueue_handler_with_options(
2978 "echo-workflow",
2979 TriggerKind::Api,
2980 json!({}),
2981 EnqueueOptions {
2982 concurrency_limits,
2983 ..Default::default()
2984 },
2985 )
2986 .await
2987 .unwrap_err();
2988 assert!(
2989 matches!(err, EngineError::InvalidConcurrencyLimit(_)),
2990 "{err:?}"
2991 );
2992 }
2993
2994 let err = engine
2996 .enqueue_handler_with_options(
2997 "not-registered",
2998 TriggerKind::Api,
2999 json!({}),
3000 EnqueueOptions {
3001 concurrency_limits: vec![ConcurrencyLimit::new("repo:acme", 0)],
3002 ..Default::default()
3003 },
3004 )
3005 .await
3006 .unwrap_err();
3007 assert!(
3008 matches!(err, EngineError::InvalidConcurrencyLimit(_)),
3009 "{err:?}"
3010 );
3011
3012 let page = engine
3013 .store()
3014 .list_runs(RunFilter::default(), 1, 10)
3015 .await
3016 .unwrap();
3017 assert_eq!(page.total, 0, "no run may be created");
3018 }
3019
3020 #[tokio::test]
3021 async fn run_handler_leaves_the_run_unattributed() {
3022 let mut engine = create_test_engine();
3023 engine.register(EchoWorkflow).unwrap();
3024
3025 let run = engine
3026 .run_handler("echo-workflow", TriggerKind::Manual, json!({}))
3027 .await
3028 .unwrap()
3029 .run;
3030
3031 assert!(run.created_by.is_none());
3032 }
3033
3034 #[tokio::test]
3035 async fn engine_register_boxed() {
3036 let mut engine = create_test_engine();
3037 let handler: Box<dyn WorkflowHandler> = Box::new(EchoWorkflow);
3038 let result = engine.register_boxed(handler);
3039 assert!(result.is_ok());
3040 assert_eq!(engine.handler_names().len(), 1);
3041 }
3042
3043 #[tokio::test]
3044 async fn engine_store_and_provider_accessors() {
3045 let store = Arc::new(InMemoryStore::new());
3046 let inner = ClaudeCodeProvider::new();
3047 let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
3048 inner,
3049 "/tmp/ironflow-fixtures",
3050 ));
3051 let engine = Engine::new(store.clone(), provider.clone());
3052
3053 let _ = engine.store();
3055 let _ = engine.provider();
3056 }
3057
3058 use crate::operation::{Operation, OperationContext};
3063 use async_trait::async_trait;
3064 use ironflow_core::error::OperationError;
3065 use ironflow_store::models::StepKind;
3066
3067 struct FakeGitlabOp {
3068 project_id: u64,
3069 title: String,
3070 }
3071
3072 #[async_trait]
3073 impl Operation for FakeGitlabOp {
3074 fn kind(&self) -> &str {
3075 "gitlab"
3076 }
3077
3078 async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
3079 Ok(json!({
3080 "issue_id": 42,
3081 "project_id": self.project_id,
3082 "title": self.title,
3083 }))
3084 }
3085
3086 fn input(&self) -> Option<Value> {
3087 Some(json!({
3088 "project_id": self.project_id,
3089 "title": self.title,
3090 }))
3091 }
3092 }
3093
3094 struct FailingOp;
3095
3096 #[async_trait]
3097 impl Operation for FailingOp {
3098 fn kind(&self) -> &str {
3099 "broken-service"
3100 }
3101
3102 async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
3103 Err(OperationError::Http {
3104 status: None,
3105 message: "service unavailable".to_string(),
3106 })
3107 }
3108 }
3109
3110 struct OperationWorkflow;
3111
3112 impl WorkflowHandler for OperationWorkflow {
3113 fn name(&self) -> &str {
3114 "operation-workflow"
3115 }
3116
3117 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3118 Box::pin(async move {
3119 let op = FakeGitlabOp {
3120 project_id: 123,
3121 title: "Bug report".to_string(),
3122 };
3123 ctx.operation("create-issue", &op).await?;
3124 Ok(())
3125 })
3126 }
3127 }
3128
3129 struct FailingOperationWorkflow;
3130
3131 impl WorkflowHandler for FailingOperationWorkflow {
3132 fn name(&self) -> &str {
3133 "failing-operation-workflow"
3134 }
3135
3136 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3137 Box::pin(async move {
3138 ctx.operation("broken-call", &FailingOp).await?;
3139 Ok(())
3140 })
3141 }
3142 }
3143
3144 struct MixedWorkflow;
3145
3146 impl WorkflowHandler for MixedWorkflow {
3147 fn name(&self) -> &str {
3148 "mixed-workflow"
3149 }
3150
3151 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3152 Box::pin(async move {
3153 ctx.shell("build", ShellConfig::new("echo built")).await?;
3154 let op = FakeGitlabOp {
3155 project_id: 456,
3156 title: "Deploy done".to_string(),
3157 };
3158 let result = ctx.operation("notify-gitlab", &op).await?;
3159 assert_eq!(result.output["issue_id"], 42);
3160 Ok(())
3161 })
3162 }
3163 }
3164
3165 #[tokio::test]
3166 async fn operation_step_happy_path() {
3167 let mut engine = create_test_engine();
3168 engine.register(OperationWorkflow).unwrap();
3169
3170 let run = engine
3171 .run_handler("operation-workflow", TriggerKind::Manual, json!({}))
3172 .await
3173 .unwrap()
3174 .run;
3175
3176 assert_eq!(run.status.state, RunStatus::Completed);
3177
3178 let steps = engine.store().list_steps(run.id).await.unwrap();
3179
3180 assert_eq!(steps.len(), 1);
3181 assert_eq!(steps[0].name, "create-issue");
3182 assert_eq!(steps[0].kind, StepKind::Custom("gitlab".to_string()));
3183 assert_eq!(
3184 steps[0].status.state,
3185 ironflow_store::models::StepStatus::Completed
3186 );
3187
3188 let output = steps[0].output.as_ref().unwrap();
3189 assert_eq!(output["issue_id"], 42);
3190 assert_eq!(output["project_id"], 123);
3191
3192 let input = steps[0].input.as_ref().unwrap();
3193 assert_eq!(input["project_id"], 123);
3194 assert_eq!(input["title"], "Bug report");
3195 }
3196
3197 #[tokio::test]
3198 async fn operation_step_failure_marks_run_failed() {
3199 let mut engine = create_test_engine();
3200 engine.register(FailingOperationWorkflow).unwrap();
3201
3202 let result = engine
3203 .run_handler("failing-operation-workflow", TriggerKind::Manual, json!({}))
3204 .await;
3205
3206 assert!(result.is_err());
3207 }
3208
3209 #[tokio::test]
3210 async fn operation_mixed_with_shell_steps() {
3211 let mut engine = create_test_engine();
3212 engine.register(MixedWorkflow).unwrap();
3213
3214 let run = engine
3215 .run_handler("mixed-workflow", TriggerKind::Manual, json!({}))
3216 .await
3217 .unwrap()
3218 .run;
3219
3220 assert_eq!(run.status.state, RunStatus::Completed);
3221
3222 let steps = engine.store().list_steps(run.id).await.unwrap();
3223
3224 assert_eq!(steps.len(), 2);
3225 assert_eq!(steps[0].kind, StepKind::Shell);
3226 assert_eq!(steps[1].kind, StepKind::Custom("gitlab".to_string()));
3227 assert_eq!(steps[0].position, 0);
3228 assert_eq!(steps[1].position, 1);
3229 }
3230
3231 use crate::config::ApprovalConfig;
3236
3237 struct SingleApprovalWorkflow;
3238
3239 impl WorkflowHandler for SingleApprovalWorkflow {
3240 fn name(&self) -> &str {
3241 "single-approval"
3242 }
3243
3244 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3245 Box::pin(async move {
3246 ctx.shell("build", ShellConfig::new("echo built")).await?;
3247 ctx.approval("gate", ApprovalConfig::new("OK?")).await?;
3248 ctx.shell("deploy", ShellConfig::new("echo deployed"))
3249 .await?;
3250 Ok(())
3251 })
3252 }
3253 }
3254
3255 struct DoubleApprovalWorkflow;
3256
3257 impl WorkflowHandler for DoubleApprovalWorkflow {
3258 fn name(&self) -> &str {
3259 "double-approval"
3260 }
3261
3262 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3263 Box::pin(async move {
3264 ctx.shell("build", ShellConfig::new("echo built")).await?;
3265 ctx.approval("staging-gate", ApprovalConfig::new("Deploy staging?"))
3266 .await?;
3267 ctx.shell("deploy-staging", ShellConfig::new("echo staging"))
3268 .await?;
3269 ctx.approval("prod-gate", ApprovalConfig::new("Deploy prod?"))
3270 .await?;
3271 ctx.shell("deploy-prod", ShellConfig::new("echo prod"))
3272 .await?;
3273 Ok(())
3274 })
3275 }
3276 }
3277
3278 #[tokio::test]
3279 async fn approval_pauses_run() {
3280 let mut engine = create_test_engine();
3281 engine.register(SingleApprovalWorkflow).unwrap();
3282
3283 let run = engine
3284 .run_handler("single-approval", TriggerKind::Manual, json!({}))
3285 .await
3286 .unwrap()
3287 .run;
3288
3289 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3290
3291 let steps = engine.store().list_steps(run.id).await.unwrap();
3292 assert_eq!(steps.len(), 2); assert_eq!(steps[0].kind, StepKind::Shell);
3294 assert_eq!(steps[0].status.state, StepStatus::Completed);
3295 assert_eq!(steps[1].kind, StepKind::Approval);
3296 assert_eq!(steps[1].status.state, StepStatus::AwaitingApproval);
3297 }
3298
3299 #[tokio::test]
3300 async fn approval_resume_completes_run() {
3301 let mut engine = create_test_engine();
3302 engine.register(SingleApprovalWorkflow).unwrap();
3303
3304 let run = engine
3306 .run_handler("single-approval", TriggerKind::Manual, json!({}))
3307 .await
3308 .unwrap()
3309 .run;
3310 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3311
3312 engine
3314 .store()
3315 .update_run_status(run.id, RunStatus::Running)
3316 .await
3317 .unwrap();
3318
3319 let resumed = engine.resume_run(run.id).await.unwrap().run;
3321 assert_eq!(resumed.status.state, RunStatus::Completed);
3322
3323 let steps = engine.store().list_steps(run.id).await.unwrap();
3324 assert_eq!(steps.len(), 3); assert_eq!(steps[0].name, "build");
3326 assert_eq!(steps[0].status.state, StepStatus::Completed);
3327 assert_eq!(steps[1].name, "gate");
3328 assert_eq!(steps[1].kind, StepKind::Approval);
3329 assert_eq!(steps[1].status.state, StepStatus::Completed);
3330 assert_eq!(steps[2].name, "deploy");
3331 assert_eq!(steps[2].status.state, StepStatus::Completed);
3332 }
3333
3334 #[tokio::test]
3335 async fn double_approval_two_resumes() {
3336 let mut engine = create_test_engine();
3337 engine.register(DoubleApprovalWorkflow).unwrap();
3338
3339 let run = engine
3341 .run_handler("double-approval", TriggerKind::Manual, json!({}))
3342 .await
3343 .unwrap()
3344 .run;
3345 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3346
3347 let steps = engine.store().list_steps(run.id).await.unwrap();
3348 assert_eq!(steps.len(), 2); engine
3352 .store()
3353 .update_run_status(run.id, RunStatus::Running)
3354 .await
3355 .unwrap();
3356
3357 let resumed = engine.resume_run(run.id).await.unwrap().run;
3358 assert_eq!(resumed.status.state, RunStatus::AwaitingApproval);
3359
3360 let steps = engine.store().list_steps(run.id).await.unwrap();
3361 assert_eq!(steps.len(), 4); engine
3365 .store()
3366 .update_run_status(run.id, RunStatus::Running)
3367 .await
3368 .unwrap();
3369
3370 let final_run = engine.resume_run(run.id).await.unwrap().run;
3371 assert_eq!(final_run.status.state, RunStatus::Completed);
3372
3373 let steps = engine.store().list_steps(run.id).await.unwrap();
3374 assert_eq!(steps.len(), 5);
3375 assert_eq!(steps[0].name, "build");
3376 assert_eq!(steps[1].name, "staging-gate");
3377 assert_eq!(steps[2].name, "deploy-staging");
3378 assert_eq!(steps[3].name, "prod-gate");
3379 assert_eq!(steps[4].name, "deploy-prod");
3380
3381 for step in &steps {
3382 assert_eq!(step.status.state, StepStatus::Completed);
3383 }
3384 }
3385
3386 use ironflow_store::models::{NewStep, StepUpdate, step_trace_id};
3391
3392 async fn create_step_with_status(
3393 store: &Arc<dyn Store>,
3394 run_id: Uuid,
3395 name: &str,
3396 position: u32,
3397 status: StepStatus,
3398 ) -> ironflow_store::models::Step {
3399 let step = store
3400 .create_step(NewStep {
3401 run_id,
3402 trace_id: step_trace_id(run_id, name, position),
3403 name: name.to_string(),
3404 kind: StepKind::Shell,
3405 position,
3406 input: None,
3407 is_error_handler: false,
3408 })
3409 .await
3410 .unwrap();
3411
3412 match status {
3413 StepStatus::Pending => {}
3414 StepStatus::Running => {
3415 store
3416 .update_step(
3417 step.id,
3418 StepUpdate {
3419 status: Some(StepStatus::Running),
3420 ..StepUpdate::default()
3421 },
3422 )
3423 .await
3424 .unwrap();
3425 }
3426 StepStatus::Completed => {
3427 store
3428 .update_step(
3429 step.id,
3430 StepUpdate {
3431 status: Some(StepStatus::Running),
3432 ..StepUpdate::default()
3433 },
3434 )
3435 .await
3436 .unwrap();
3437 store
3438 .update_step(
3439 step.id,
3440 StepUpdate {
3441 status: Some(StepStatus::Completed),
3442 ..StepUpdate::default()
3443 },
3444 )
3445 .await
3446 .unwrap();
3447 }
3448 StepStatus::AwaitingApproval => {
3449 store
3450 .update_step(
3451 step.id,
3452 StepUpdate {
3453 status: Some(StepStatus::Running),
3454 ..StepUpdate::default()
3455 },
3456 )
3457 .await
3458 .unwrap();
3459 store
3460 .update_step(
3461 step.id,
3462 StepUpdate {
3463 status: Some(StepStatus::AwaitingApproval),
3464 ..StepUpdate::default()
3465 },
3466 )
3467 .await
3468 .unwrap();
3469 }
3470 _ => panic!("unsupported status for test helper: {status}"),
3471 }
3472
3473 store.get_step(step.id).await.unwrap().unwrap()
3474 }
3475
3476 #[tokio::test]
3477 async fn fail_orphaned_steps_marks_running_as_failed() {
3478 let engine = create_test_engine();
3479 let run = engine
3480 .store()
3481 .create_run(NewRun {
3482 created_by: None,
3483 workflow_name: "test".to_string(),
3484 trigger: TriggerKind::Manual,
3485 payload: json!({}),
3486 max_retries: 0,
3487 handler_version: None,
3488 labels: HashMap::new(),
3489 scheduled_at: None,
3490 idempotency_key: None,
3491 concurrency_key: None,
3492 concurrency_limits: Vec::new(),
3493 max_cost_usd: None,
3494 })
3495 .await
3496 .unwrap()
3497 .into_run();
3498
3499 let step = create_step_with_status(
3500 engine.store(),
3501 run.id,
3502 "running-step",
3503 0,
3504 StepStatus::Running,
3505 )
3506 .await;
3507
3508 engine
3509 .fail_orphaned_steps(run.id, "parent run timed out")
3510 .await
3511 .unwrap();
3512
3513 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
3514 assert_eq!(updated.status.state, StepStatus::Failed);
3515 assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
3516 assert!(updated.completed_at.is_some());
3517 }
3518
3519 #[tokio::test]
3520 async fn fail_orphaned_steps_marks_pending_as_skipped() {
3521 let engine = create_test_engine();
3522 let run = engine
3523 .store()
3524 .create_run(NewRun {
3525 created_by: None,
3526 workflow_name: "test".to_string(),
3527 trigger: TriggerKind::Manual,
3528 payload: json!({}),
3529 max_retries: 0,
3530 handler_version: None,
3531 labels: HashMap::new(),
3532 scheduled_at: None,
3533 idempotency_key: None,
3534 concurrency_key: None,
3535 concurrency_limits: Vec::new(),
3536 max_cost_usd: None,
3537 })
3538 .await
3539 .unwrap()
3540 .into_run();
3541
3542 let step = create_step_with_status(
3543 engine.store(),
3544 run.id,
3545 "pending-step",
3546 0,
3547 StepStatus::Pending,
3548 )
3549 .await;
3550
3551 engine
3552 .fail_orphaned_steps(run.id, "parent run timed out")
3553 .await
3554 .unwrap();
3555
3556 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
3557 assert_eq!(updated.status.state, StepStatus::Skipped);
3558 assert!(updated.error.is_none());
3559 assert!(updated.completed_at.is_some());
3560 }
3561
3562 #[tokio::test]
3563 async fn fail_orphaned_steps_marks_awaiting_approval_as_failed() {
3564 let engine = create_test_engine();
3565 let run = engine
3566 .store()
3567 .create_run(NewRun {
3568 created_by: None,
3569 workflow_name: "test".to_string(),
3570 trigger: TriggerKind::Manual,
3571 payload: json!({}),
3572 max_retries: 0,
3573 handler_version: None,
3574 labels: HashMap::new(),
3575 scheduled_at: None,
3576 idempotency_key: None,
3577 concurrency_key: None,
3578 concurrency_limits: Vec::new(),
3579 max_cost_usd: None,
3580 })
3581 .await
3582 .unwrap()
3583 .into_run();
3584
3585 let step = create_step_with_status(
3586 engine.store(),
3587 run.id,
3588 "approval-step",
3589 0,
3590 StepStatus::AwaitingApproval,
3591 )
3592 .await;
3593
3594 engine
3595 .fail_orphaned_steps(run.id, "parent run timed out")
3596 .await
3597 .unwrap();
3598
3599 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
3600 assert_eq!(updated.status.state, StepStatus::Failed);
3601 assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
3602 assert!(updated.completed_at.is_some());
3603 }
3604
3605 #[tokio::test]
3606 async fn fail_orphaned_steps_skips_terminal_steps() {
3607 let engine = create_test_engine();
3608 let run = engine
3609 .store()
3610 .create_run(NewRun {
3611 created_by: None,
3612 workflow_name: "test".to_string(),
3613 trigger: TriggerKind::Manual,
3614 payload: json!({}),
3615 max_retries: 0,
3616 handler_version: None,
3617 labels: HashMap::new(),
3618 scheduled_at: None,
3619 idempotency_key: None,
3620 concurrency_key: None,
3621 concurrency_limits: Vec::new(),
3622 max_cost_usd: None,
3623 })
3624 .await
3625 .unwrap()
3626 .into_run();
3627
3628 let completed_step =
3629 create_step_with_status(engine.store(), run.id, "done", 0, StepStatus::Completed).await;
3630 let running_step =
3631 create_step_with_status(engine.store(), run.id, "in-flight", 1, StepStatus::Running)
3632 .await;
3633
3634 engine
3635 .fail_orphaned_steps(run.id, "parent run timed out")
3636 .await
3637 .unwrap();
3638
3639 let completed = engine
3640 .store()
3641 .get_step(completed_step.id)
3642 .await
3643 .unwrap()
3644 .unwrap();
3645 assert_eq!(completed.status.state, StepStatus::Completed);
3646
3647 let failed = engine
3648 .store()
3649 .get_step(running_step.id)
3650 .await
3651 .unwrap()
3652 .unwrap();
3653 assert_eq!(failed.status.state, StepStatus::Failed);
3654 }
3655
3656 #[tokio::test]
3657 async fn fail_orphaned_steps_mixed_states() {
3658 let engine = create_test_engine();
3659 let run = engine
3660 .store()
3661 .create_run(NewRun {
3662 created_by: None,
3663 workflow_name: "test".to_string(),
3664 trigger: TriggerKind::Manual,
3665 payload: json!({}),
3666 max_retries: 0,
3667 handler_version: None,
3668 labels: HashMap::new(),
3669 scheduled_at: None,
3670 idempotency_key: None,
3671 concurrency_key: None,
3672 concurrency_limits: Vec::new(),
3673 max_cost_usd: None,
3674 })
3675 .await
3676 .unwrap()
3677 .into_run();
3678
3679 let s_completed =
3680 create_step_with_status(engine.store(), run.id, "step-1", 0, StepStatus::Completed)
3681 .await;
3682 let s_running =
3683 create_step_with_status(engine.store(), run.id, "step-2", 1, StepStatus::Running).await;
3684 let s_pending =
3685 create_step_with_status(engine.store(), run.id, "step-3", 2, StepStatus::Pending).await;
3686
3687 engine.fail_orphaned_steps(run.id, "timeout").await.unwrap();
3688
3689 let r_completed = engine
3690 .store()
3691 .get_step(s_completed.id)
3692 .await
3693 .unwrap()
3694 .unwrap();
3695 assert_eq!(r_completed.status.state, StepStatus::Completed);
3696
3697 let r_running = engine
3698 .store()
3699 .get_step(s_running.id)
3700 .await
3701 .unwrap()
3702 .unwrap();
3703 assert_eq!(r_running.status.state, StepStatus::Failed);
3704 assert_eq!(r_running.error.as_deref(), Some("timeout"));
3705
3706 let r_pending = engine
3707 .store()
3708 .get_step(s_pending.id)
3709 .await
3710 .unwrap()
3711 .unwrap();
3712 assert_eq!(r_pending.status.state, StepStatus::Skipped);
3713 assert!(r_pending.error.is_none());
3714 }
3715
3716 #[tokio::test]
3717 async fn fail_orphaned_steps_no_steps_is_noop() {
3718 let engine = create_test_engine();
3719 let run = engine
3720 .store()
3721 .create_run(NewRun {
3722 created_by: None,
3723 workflow_name: "test".to_string(),
3724 trigger: TriggerKind::Manual,
3725 payload: json!({}),
3726 max_retries: 0,
3727 handler_version: None,
3728 labels: HashMap::new(),
3729 scheduled_at: None,
3730 idempotency_key: None,
3731 concurrency_key: None,
3732 concurrency_limits: Vec::new(),
3733 max_cost_usd: None,
3734 })
3735 .await
3736 .unwrap()
3737 .into_run();
3738
3739 let result = engine.fail_orphaned_steps(run.id, "timeout").await;
3740 assert!(result.is_ok());
3741 }
3742
3743 #[tokio::test]
3744 async fn fail_orphaned_steps_preserves_existing_error() {
3745 let engine = create_test_engine();
3746 let run = engine
3747 .store()
3748 .create_run(NewRun {
3749 created_by: None,
3750 workflow_name: "test".to_string(),
3751 trigger: TriggerKind::Manual,
3752 payload: json!({}),
3753 max_retries: 0,
3754 handler_version: None,
3755 labels: HashMap::new(),
3756 scheduled_at: None,
3757 idempotency_key: None,
3758 concurrency_key: None,
3759 concurrency_limits: Vec::new(),
3760 max_cost_usd: None,
3761 })
3762 .await
3763 .unwrap()
3764 .into_run();
3765
3766 let step_with_error = create_step_with_status(
3767 engine.store(),
3768 run.id,
3769 "already-errored",
3770 0,
3771 StepStatus::Running,
3772 )
3773 .await;
3774
3775 engine
3776 .store()
3777 .update_step(
3778 step_with_error.id,
3779 StepUpdate {
3780 error: Some("real error from provider".to_string()),
3781 ..StepUpdate::default()
3782 },
3783 )
3784 .await
3785 .unwrap();
3786
3787 let step_no_error = create_step_with_status(
3788 engine.store(),
3789 run.id,
3790 "no-error-yet",
3791 1,
3792 StepStatus::Running,
3793 )
3794 .await;
3795
3796 engine
3797 .fail_orphaned_steps(run.id, "parent run failed")
3798 .await
3799 .unwrap();
3800
3801 let updated_with = engine
3802 .store()
3803 .get_step(step_with_error.id)
3804 .await
3805 .unwrap()
3806 .unwrap();
3807 assert_eq!(updated_with.status.state, StepStatus::Failed);
3808 assert_eq!(
3809 updated_with.error.as_deref(),
3810 Some("real error from provider"),
3811 );
3812
3813 let updated_without = engine
3814 .store()
3815 .get_step(step_no_error.id)
3816 .await
3817 .unwrap()
3818 .unwrap();
3819 assert_eq!(updated_without.status.state, StepStatus::Failed);
3820 assert_eq!(updated_without.error.as_deref(), Some("parent run failed"),);
3821 }
3822}