1use std::collections::HashMap;
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;
17use tracing::{error, info, warn};
18use uuid::Uuid;
19
20#[cfg(feature = "prometheus")]
21use ironflow_core::metric_names::{
22 RUN_BUDGET_EXCEEDED_TOTAL, RUN_COST_USD, RUN_DURATION_SECONDS, RUNS_ACTIVE, RUNS_TOTAL,
23};
24use ironflow_core::provider::AgentProvider;
25use ironflow_store::error::StoreError;
26use ironflow_store::models::{
27 NewRun, Run, RunActor, RunCreation, RunFilter, RunStatus, RunUpdate, StepStatus, StepUpdate,
28 TriggerKind,
29};
30use ironflow_store::store::Store;
31#[cfg(feature = "prometheus")]
32use metrics::{counter, gauge, histogram};
33
34use crate::artifact::ArtifactSink;
35use crate::budget::{BudgetConfig, month_start};
36use crate::context::WorkflowContext;
37use crate::error::EngineError;
38use crate::executor::StepResult;
39use crate::guard::{WorkflowGuardConfig, new_shared_guard_state};
40use crate::handler::{WorkflowHandler, WorkflowInfo};
41use crate::log_sender::LogSender;
42use crate::notify::{
43 Event, EventPublisher, EventSubscriber, RunBudgetExceededEvent, RunFailedEvent,
44 RunStatusChangedEvent, WorkflowEventBus,
45};
46use crate::plan::{
47 ExecutionPlan, PlanOptions, PlanRecorder, SharedPlanRecorder, estimate_durations, lock_plan,
48};
49use crate::retry_policy::{backoff_for_retry, is_run_retryable};
50use crate::schedule::CronSchedule;
51use ironflow_core::decision::DecisionProvider;
52
53#[derive(Debug, Clone)]
72pub struct WorkflowResult {
73 pub run: Run,
75 pub steps: Vec<StepResult>,
77}
78
79#[derive(Debug, Clone, Default)]
98pub struct EnqueueOptions {
99 pub max_retries: u32,
101 pub labels: HashMap<String, String>,
103 pub scheduled_at: Option<DateTime<Utc>>,
106 pub max_cost_usd: Option<Decimal>,
110 pub created_by: Option<RunActor>,
113 pub idempotency_key: Option<String>,
119}
120
121pub struct Engine {
162 store: Arc<dyn Store>,
163 provider: Arc<dyn AgentProvider>,
164 handlers: HashMap<String, Arc<dyn WorkflowHandler>>,
165 event_publisher: EventPublisher,
166 log_sender: Option<LogSender>,
167 budget: BudgetConfig,
168 artifact_sink: Option<Arc<dyn ArtifactSink>>,
169 guard_config: Option<WorkflowGuardConfig>,
170 event_bus: Option<WorkflowEventBus>,
171 decision_provider: Option<Arc<dyn DecisionProvider>>,
172}
173
174fn validate_category(handler_name: &str, category: &str) -> Result<(), EngineError> {
184 let reject = |reason: &str| {
185 Err(EngineError::InvalidWorkflow(format!(
186 "handler '{handler_name}' has invalid category '{category}': {reason}"
187 )))
188 };
189
190 if category.is_empty() {
191 return reject("empty category");
192 }
193 if category.starts_with('/') {
194 return reject("leading '/'");
195 }
196 if category.ends_with('/') {
197 return reject("trailing '/'");
198 }
199 for segment in category.split('/') {
200 if segment.is_empty() {
201 return reject("empty segment (double '/')");
202 }
203 if segment.trim().is_empty() {
204 return reject("whitespace-only segment");
205 }
206 }
207 Ok(())
208}
209
210impl Engine {
211 pub fn new(store: Arc<dyn Store>, provider: Arc<dyn AgentProvider>) -> Self {
227 Self {
228 store,
229 provider,
230 handlers: HashMap::new(),
231 event_publisher: EventPublisher::new(),
232 log_sender: None,
233 budget: BudgetConfig::new(),
234 artifact_sink: None,
235 guard_config: None,
236 event_bus: None,
237 decision_provider: None,
238 }
239 }
240
241 pub fn with_decision_provider(mut self, provider: Arc<dyn DecisionProvider>) -> Self {
262 self.decision_provider = Some(provider);
263 self
264 }
265
266 pub fn with_budget_config(mut self, budget: BudgetConfig) -> Self {
287 self.budget = budget;
288 self
289 }
290
291 pub fn budget_config(&self) -> &BudgetConfig {
293 &self.budget
294 }
295
296 pub fn with_guard_config(mut self, config: WorkflowGuardConfig) -> Self {
318 self.guard_config = Some(config);
319 self
320 }
321
322 pub fn guard_config(&self) -> Option<&WorkflowGuardConfig> {
324 self.guard_config.as_ref()
325 }
326
327 pub fn set_log_sender(&mut self, sender: LogSender) {
333 self.log_sender = Some(sender);
334 }
335
336 pub fn set_artifact_sink(&mut self, sink: Arc<dyn ArtifactSink>) {
354 self.artifact_sink = Some(sink);
355 }
356
357 pub fn artifact_sink(&self) -> Option<&Arc<dyn ArtifactSink>> {
359 self.artifact_sink.as_ref()
360 }
361
362 pub fn set_event_bus(&mut self, bus: WorkflowEventBus) {
379 self.event_bus = Some(bus);
380 }
381
382 pub fn event_bus(&self) -> Option<&WorkflowEventBus> {
384 self.event_bus.as_ref()
385 }
386
387 pub fn store(&self) -> &Arc<dyn Store> {
389 &self.store
390 }
391
392 pub fn provider(&self) -> &Arc<dyn AgentProvider> {
394 &self.provider
395 }
396
397 fn build_context(&self, run: &Run) -> WorkflowContext {
406 let handlers = self.handlers.clone();
407 let resolver: crate::context::HandlerResolver =
408 Arc::new(move |name: &str| handlers.get(name).cloned());
409 let mut ctx = WorkflowContext::with_handler_resolver(
410 run.id,
411 run.workflow_name.clone(),
412 self.store.clone(),
413 self.provider.clone(),
414 resolver,
415 );
416 ctx.carry_over_run_totals(run.retry_count + 1, run.cost_usd, run.duration_ms);
417 ctx.set_max_cost_usd(run.max_cost_usd);
418 if let Some(ref sender) = self.log_sender {
419 ctx.set_log_sender(sender.clone());
420 }
421 if let Some(ref sink) = self.artifact_sink {
422 ctx.set_artifact_sink(sink.clone());
423 }
424 if let Some(ref bus) = self.event_bus {
425 ctx.set_event_bus(bus.clone());
426 }
427 if let Some(ref provider) = self.decision_provider {
428 ctx.set_decision_provider(provider.clone());
429 }
430 ctx
431 }
432
433 fn build_context_with_guard(
439 &self,
440 run: &Run,
441 handler: &dyn WorkflowHandler,
442 ) -> WorkflowContext {
443 let mut ctx = self.build_context(run);
444 let guard_config = handler.guard_config().or_else(|| self.guard_config.clone());
445 if let Some(config) = guard_config {
446 ctx.set_guard(config, new_shared_guard_state());
447 }
448 ctx
449 }
450
451 async fn check_monthly_quota(&self, workflow_name: &str) -> Result<(), EngineError> {
462 let Some(limit) = self.budget.monthly_cost_limit_usd else {
463 return Ok(());
464 };
465
466 let stats = self
467 .store
468 .get_stats(RunFilter {
469 created_after: Some(month_start(Utc::now())),
470 ..RunFilter::default()
471 })
472 .await?;
473
474 if stats.total_cost_usd < limit {
475 return Ok(());
476 }
477
478 warn!(
479 workflow = %workflow_name,
480 limit_usd = %limit,
481 spent_usd = %stats.total_cost_usd,
482 "monthly cost quota exhausted, refusing new run"
483 );
484
485 #[cfg(feature = "prometheus")]
486 counter!(
487 RUN_BUDGET_EXCEEDED_TOTAL,
488 "workflow" => workflow_name.to_string(),
489 "scope" => "monthly",
490 )
491 .increment(1);
492
493 Err(EngineError::MonthlyBudgetExceeded {
494 limit_usd: limit,
495 spent_usd: stats.total_cost_usd,
496 })
497 }
498
499 pub fn register(&mut self, handler: impl WorkflowHandler + 'static) -> Result<(), EngineError> {
543 let name = handler.name().to_string();
544 if self.handlers.contains_key(&name) {
545 return Err(EngineError::InvalidWorkflow(format!(
546 "handler '{}' already registered",
547 name
548 )));
549 }
550 if let Some(category) = handler.category() {
551 validate_category(&name, category)?;
552 }
553 self.handlers.insert(name, Arc::new(handler));
554 Ok(())
555 }
556
557 pub fn register_boxed(&mut self, handler: Box<dyn WorkflowHandler>) -> Result<(), EngineError> {
564 let name = handler.name().to_string();
565 if self.handlers.contains_key(&name) {
566 return Err(EngineError::InvalidWorkflow(format!(
567 "handler '{}' already registered",
568 name
569 )));
570 }
571 if let Some(category) = handler.category() {
572 validate_category(&name, category)?;
573 }
574 self.handlers.insert(name, Arc::from(handler));
575 Ok(())
576 }
577
578 pub fn get_handler(&self, name: &str) -> Option<&Arc<dyn WorkflowHandler>> {
580 self.handlers.get(name)
581 }
582
583 pub fn handler_names(&self) -> Vec<&str> {
585 self.handlers.keys().map(|s| s.as_str()).collect()
586 }
587
588 pub fn handler_info(&self, name: &str) -> Option<WorkflowInfo> {
590 self.handlers.get(name).map(|h| h.describe())
591 }
592
593 pub fn scheduled_handlers(&self) -> Vec<(&str, &CronSchedule)> {
617 self.handlers
618 .iter()
619 .filter_map(|(name, handler)| handler.schedule().map(|sched| (name.as_str(), sched)))
620 .collect()
621 }
622
623 pub fn subscribe(
648 &mut self,
649 subscriber: impl EventSubscriber + 'static,
650 event_types: &[&'static str],
651 ) {
652 self.event_publisher.subscribe(subscriber, event_types);
653 }
654
655 pub fn event_publisher(&self) -> &EventPublisher {
660 &self.event_publisher
661 }
662
663 #[tracing::instrument(name = "engine.run_handler", skip_all, fields(workflow = %handler_name))]
693 pub async fn run_handler(
694 &self,
695 handler_name: &str,
696 trigger: TriggerKind,
697 payload: Value,
698 ) -> Result<WorkflowResult, EngineError> {
699 let handler = self
700 .handlers
701 .get(handler_name)
702 .ok_or_else(|| {
703 EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
704 })?
705 .clone();
706
707 self.check_monthly_quota(handler_name).await?;
708
709 let handler_version = handler.version().map(str::to_string);
710 let max_cost_usd = self
711 .budget
712 .resolve_run_cap(None, handler.default_max_cost_usd());
713 let run = self
714 .store
715 .create_run(NewRun {
716 created_by: None,
717 workflow_name: handler_name.to_string(),
718 trigger,
719 payload,
720 max_retries: 0,
721 handler_version,
722 labels: handler.default_labels(),
723 scheduled_at: None,
724 idempotency_key: None,
725 max_cost_usd,
726 })
727 .await?
728 .into_run();
729
730 let run_id = run.id;
731 info!(run_id = %run_id, handler_version = run.handler_version.as_deref().unwrap_or(""), "run created");
732
733 self.store
734 .update_run_status(run_id, RunStatus::Running)
735 .await?;
736
737 #[cfg(feature = "prometheus")]
738 gauge!(RUNS_ACTIVE, "workflow" => handler_name.to_string()).increment(1.0);
739
740 let run_start = Instant::now();
741 let mut ctx = self.build_context_with_guard(&run, handler.as_ref());
742
743 let result = handler.execute(&mut ctx).await;
744 self.finalize_run(run_id, handler_name, result, &ctx, run_start, run.labels)
745 .await
746 }
747
748 #[tracing::instrument(name = "engine.plan_handler", skip_all, fields(workflow = %handler_name))]
790 pub async fn plan_handler(
791 &self,
792 handler_name: &str,
793 payload: Value,
794 options: PlanOptions,
795 ) -> Result<ExecutionPlan, EngineError> {
796 if options.max_depth == 0 {
797 return Err(EngineError::InvalidWorkflow(
798 "max_depth must be at least 1".to_string(),
799 ));
800 }
801
802 let handler = self
803 .handlers
804 .get(handler_name)
805 .ok_or_else(|| {
806 EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
807 })?
808 .clone();
809
810 let estimates = if options.estimate_durations {
811 estimate_durations(&self.store, handler_name, options.sample_runs).await?
812 } else {
813 HashMap::new()
814 };
815
816 let shared: SharedPlanRecorder = Arc::new(Mutex::new(PlanRecorder::new(
817 handler_name.to_string(),
818 payload,
819 options.max_depth,
820 estimates,
821 )));
822
823 let handlers = self.handlers.clone();
826 let resolver: crate::context::HandlerResolver =
827 Arc::new(move |name: &str| handlers.get(name).cloned());
828 let mut ctx = WorkflowContext::with_handler_resolver(
829 Uuid::now_v7(),
830 handler_name.to_string(),
831 self.store.clone(),
832 self.provider.clone(),
833 resolver,
834 );
835 ctx.set_plan(shared.clone());
836
837 if let Err(err) = handler.execute(&mut ctx).await {
838 lock_plan(&shared).fail(err.to_string());
839 }
840 drop(ctx);
841
842 let plan = match Arc::try_unwrap(shared) {
843 Ok(mutex) => mutex
844 .into_inner()
845 .unwrap_or_else(|poisoned| poisoned.into_inner())
846 .into_plan(),
847 Err(shared) => lock_plan(&shared).snapshot(),
848 };
849
850 info!(
851 workflow = %handler_name,
852 steps = plan.steps.len(),
853 truncated = plan.truncated,
854 "execution plan built"
855 );
856
857 Ok(plan)
858 }
859
860 #[tracing::instrument(name = "engine.enqueue_handler", skip_all, fields(workflow = %handler_name))]
871 pub async fn enqueue_handler(
872 &self,
873 handler_name: &str,
874 trigger: TriggerKind,
875 payload: Value,
876 max_retries: u32,
877 ) -> Result<Run, EngineError> {
878 self.enqueue_handler_with_options(
879 handler_name,
880 trigger,
881 payload,
882 EnqueueOptions {
883 max_retries,
884 ..Default::default()
885 },
886 )
887 .await
888 .map(RunCreation::into_run)
889 }
890
891 #[tracing::instrument(name = "engine.enqueue_handler_with_options", skip_all, fields(workflow = %handler_name))]
936 pub async fn enqueue_handler_with_options(
937 &self,
938 handler_name: &str,
939 trigger: TriggerKind,
940 payload: Value,
941 options: EnqueueOptions,
942 ) -> Result<RunCreation, EngineError> {
943 let EnqueueOptions {
944 max_retries,
945 labels,
946 scheduled_at,
947 max_cost_usd,
948 created_by,
949 idempotency_key,
950 } = options;
951
952 let handler = self.handlers.get(handler_name).ok_or_else(|| {
953 EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
954 })?;
955
956 self.check_monthly_quota(handler_name).await?;
957
958 let handler_version = handler.version().map(str::to_string);
959 let mut merged_labels = handler.default_labels();
960 merged_labels.extend(labels);
961 let resolved_cap = self
962 .budget
963 .resolve_run_cap(max_cost_usd, handler.default_max_cost_usd());
964
965 let creation = self
966 .store
967 .create_run(NewRun {
968 workflow_name: handler_name.to_string(),
969 trigger,
970 payload,
971 max_retries,
972 handler_version,
973 labels: merged_labels,
974 scheduled_at,
975 created_by,
976 idempotency_key,
977 max_cost_usd: resolved_cap,
978 })
979 .await?;
980
981 match &creation {
982 RunCreation::Created(run) => info!(
983 run_id = %run.id,
984 workflow = %handler_name,
985 max_cost_usd = ?resolved_cap,
986 "handler run enqueued"
987 ),
988 RunCreation::Existing(run) => info!(
989 run_id = %run.id,
990 workflow = %handler_name,
991 "idempotent replay, nothing enqueued"
992 ),
993 }
994
995 Ok(creation)
996 }
997
998 #[tracing::instrument(name = "engine.execute_handler_run", skip_all, fields(run_id = %run_id))]
1007 pub async fn execute_handler_run(&self, run_id: Uuid) -> Result<WorkflowResult, EngineError> {
1008 let run = self
1009 .store
1010 .get_run(run_id)
1011 .await?
1012 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
1013
1014 let handler = self
1015 .handlers
1016 .get(&run.workflow_name)
1017 .ok_or_else(|| {
1018 EngineError::InvalidWorkflow(format!(
1019 "no handler registered: {}",
1020 run.workflow_name
1021 ))
1022 })?
1023 .clone();
1024
1025 #[cfg(feature = "prometheus")]
1026 gauge!(RUNS_ACTIVE, "workflow" => run.workflow_name.clone()).increment(1.0);
1027
1028 let run_start = Instant::now();
1029 let mut ctx = self.build_context_with_guard(&run, handler.as_ref());
1030
1031 if run.retry_count > 0 {
1034 ctx.load_replay_steps().await?;
1035 }
1036
1037 let result = handler.execute(&mut ctx).await;
1038 self.finalize_run(
1039 run_id,
1040 &run.workflow_name,
1041 result,
1042 &ctx,
1043 run_start,
1044 run.labels,
1045 )
1046 .await
1047 }
1048
1049 #[tracing::instrument(name = "engine.execute_run", skip_all, fields(run_id = %run_id))]
1057 pub async fn execute_run(&self, run_id: Uuid) -> Result<WorkflowResult, EngineError> {
1058 self.execute_handler_run(run_id).await
1059 }
1060
1061 #[tracing::instrument(name = "engine.resume_run", skip_all, fields(run_id = %run_id))]
1075 pub async fn resume_run(&self, run_id: Uuid) -> Result<WorkflowResult, EngineError> {
1076 let run = self
1077 .store
1078 .get_run(run_id)
1079 .await?
1080 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
1081
1082 let handler = self
1083 .handlers
1084 .get(&run.workflow_name)
1085 .ok_or_else(|| {
1086 EngineError::InvalidWorkflow(format!(
1087 "no handler registered: {}",
1088 run.workflow_name
1089 ))
1090 })?
1091 .clone();
1092
1093 info!(run_id = %run_id, workflow = %run.workflow_name, "resuming run after approval");
1094
1095 let run_start = Instant::now();
1096 let mut ctx = self.build_context_with_guard(&run, handler.as_ref());
1097 ctx.load_replay_steps().await?;
1098
1099 let result = handler.execute(&mut ctx).await;
1100 self.finalize_run(
1101 run_id,
1102 &run.workflow_name,
1103 result,
1104 &ctx,
1105 run_start,
1106 run.labels,
1107 )
1108 .await
1109 }
1110
1111 pub async fn fail_or_schedule_retry(
1156 &self,
1157 run_id: Uuid,
1158 error: &str,
1159 retryable: bool,
1160 cost_usd: Option<Decimal>,
1161 duration_ms: Option<u64>,
1162 ) -> Result<RunStatus, EngineError> {
1163 let run = self
1164 .store
1165 .get_run(run_id)
1166 .await?
1167 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
1168
1169 let has_attempts_left = run.retry_count < run.max_retries;
1170 let update = if retryable && has_attempts_left {
1171 let backoff = backoff_for_retry(run.retry_count);
1172 let scheduled_at = Utc::now() + TimeDelta::milliseconds(backoff.as_millis() as i64);
1173
1174 info!(
1175 run_id = %run_id,
1176 workflow = %run.workflow_name,
1177 attempt = run.retry_count + 1,
1178 max_retries = run.max_retries,
1179 backoff_secs = backoff.as_secs(),
1180 scheduled_at = %scheduled_at,
1181 "run failed, scheduling retry"
1182 );
1183
1184 RunUpdate {
1185 status: Some(RunStatus::Retrying),
1186 error: Some(error.to_string()),
1187 increment_retry: true,
1188 cost_usd,
1189 duration_ms,
1190 scheduled_at: Some(scheduled_at),
1191 ..RunUpdate::default()
1192 }
1193 } else {
1194 RunUpdate {
1195 status: Some(RunStatus::Failed),
1196 error: Some(error.to_string()),
1197 cost_usd,
1198 duration_ms,
1199 completed_at: Some(Utc::now()),
1200 ..RunUpdate::default()
1201 }
1202 };
1203
1204 let status = update.status.unwrap_or(RunStatus::Failed);
1205 self.store.update_run(run_id, update).await?;
1206 self.fail_orphaned_steps(run_id, error).await?;
1207
1208 Ok(status)
1209 }
1210
1211 pub async fn fail_orphaned_steps(
1225 &self,
1226 run_id: Uuid,
1227 error_message: &str,
1228 ) -> Result<(), EngineError> {
1229 let steps = self.store.list_steps(run_id).await?;
1230 let now = Utc::now();
1231
1232 for step in steps {
1233 if step.status.state.is_terminal() {
1234 continue;
1235 }
1236
1237 let (target_status, error) = match step.status.state {
1238 StepStatus::Running | StepStatus::AwaitingApproval => {
1239 let err = if step.error.is_some() {
1240 None
1241 } else {
1242 Some(error_message.to_string())
1243 };
1244 (StepStatus::Failed, err)
1245 }
1246 StepStatus::Pending => (StepStatus::Skipped, None),
1247 _ => continue,
1248 };
1249
1250 if let Err(e) = self
1251 .store
1252 .update_step(
1253 step.id,
1254 StepUpdate {
1255 status: Some(target_status),
1256 error,
1257 completed_at: Some(now),
1258 ..StepUpdate::default()
1259 },
1260 )
1261 .await
1262 {
1263 warn!(
1264 run_id = %run_id,
1265 step_id = %step.id,
1266 step_name = %step.name,
1267 error = %e,
1268 "failed to cleanup orphaned step"
1269 );
1270 } else {
1271 info!(
1272 run_id = %run_id,
1273 step_id = %step.id,
1274 step_name = %step.name,
1275 from = %step.status.state,
1276 to = %target_status,
1277 "cleaned up orphaned step"
1278 );
1279 }
1280 }
1281
1282 Ok(())
1283 }
1284
1285 async fn finalize_run(
1291 &self,
1292 run_id: Uuid,
1293 workflow_name: &str,
1294 result: Result<(), EngineError>,
1295 ctx: &WorkflowContext,
1296 run_start: Instant,
1297 run_labels: HashMap<String, String>,
1298 ) -> Result<WorkflowResult, EngineError> {
1299 let total_duration = ctx.carried_duration_ms() + run_start.elapsed().as_millis() as u64;
1302 let completed_at = Utc::now();
1303
1304 let final_status;
1305 let final_run;
1306
1307 match result {
1308 Ok(()) => {
1309 final_status = if ctx.has_allowed_failure() {
1310 RunStatus::Warning
1311 } else {
1312 RunStatus::Completed
1313 };
1314 final_run = self
1315 .store
1316 .update_run_returning(
1317 run_id,
1318 RunUpdate {
1319 status: Some(final_status),
1320 cost_usd: Some(ctx.total_cost_usd()),
1321 duration_ms: Some(total_duration),
1322 completed_at: Some(completed_at),
1323 ..RunUpdate::default()
1324 },
1325 )
1326 .await?;
1327
1328 info!(
1329 run_id = %run_id,
1330 status = %final_status,
1331 cost_usd = %ctx.total_cost_usd(),
1332 duration_ms = total_duration,
1333 "run completed"
1334 );
1335 }
1336 Err(EngineError::ApprovalRequired {
1337 run_id: approval_run_id,
1338 step_id,
1339 ref message,
1340 }) => {
1341 final_status = RunStatus::AwaitingApproval;
1342 final_run = self
1343 .store
1344 .update_run_returning(
1345 run_id,
1346 RunUpdate {
1347 status: Some(RunStatus::AwaitingApproval),
1348 cost_usd: Some(ctx.total_cost_usd()),
1349 duration_ms: Some(total_duration),
1350 ..RunUpdate::default()
1351 },
1352 )
1353 .await?;
1354
1355 info!(
1356 run_id = %approval_run_id,
1357 step_id = %step_id,
1358 message = %message,
1359 "run awaiting approval"
1360 );
1361 }
1362 Err(EngineError::DelaySleeping {
1363 run_id: delay_run_id,
1364 step_id,
1365 wake_at,
1366 }) => {
1367 final_status = RunStatus::Sleeping;
1368 final_run = self
1369 .store
1370 .update_run_returning(
1371 run_id,
1372 RunUpdate {
1373 status: Some(RunStatus::Sleeping),
1374 cost_usd: Some(ctx.total_cost_usd()),
1375 duration_ms: Some(total_duration),
1376 scheduled_at: Some(wake_at),
1377 ..RunUpdate::default()
1378 },
1379 )
1380 .await?;
1381
1382 info!(
1383 run_id = %delay_run_id,
1384 step_id = %step_id,
1385 wake_at = %wake_at,
1386 "run sleeping until delay elapses"
1387 );
1388 }
1389 Err(err) => {
1390 let guardrail_stop = matches!(
1394 err,
1395 EngineError::RunBudgetExceeded { .. } | EngineError::WorkflowGuardRejected(_)
1396 );
1397
1398 final_status = if guardrail_stop {
1399 if let Err(store_err) = self
1400 .store
1401 .update_run(
1402 run_id,
1403 RunUpdate {
1404 status: Some(RunStatus::Cancelled),
1405 error: Some(err.to_string()),
1406 cost_usd: Some(ctx.total_cost_usd()),
1407 duration_ms: Some(total_duration),
1408 completed_at: Some(completed_at),
1409 ..RunUpdate::default()
1410 },
1411 )
1412 .await
1413 {
1414 error!(run_id = %run_id, store_error = %store_err, "failed to persist run cancellation");
1415 }
1416 if let Err(cleanup_err) = self
1417 .fail_orphaned_steps(run_id, "run stopped: guardrail limit reached")
1418 .await
1419 {
1420 error!(run_id = %run_id, store_error = %cleanup_err, "failed to cleanup orphaned steps");
1421 }
1422 RunStatus::Cancelled
1423 } else {
1424 self.fail_or_schedule_retry(
1425 run_id,
1426 &err.to_string(),
1427 is_run_retryable(&err),
1428 Some(ctx.total_cost_usd()),
1429 Some(total_duration),
1430 )
1431 .await
1432 .unwrap_or_else(|store_err| {
1433 error!(run_id = %run_id, store_error = %store_err, "failed to persist run failure");
1434 RunStatus::Failed
1435 })
1436 };
1437
1438 if matches!(err, EngineError::RunBudgetExceeded { .. }) {
1439 self.on_run_budget_exceeded(workflow_name, run_id, &err);
1440 }
1441
1442 error!(run_id = %run_id, status = %final_status, error = %err, "run stopped");
1443
1444 self.publish_run_status_changed(
1445 workflow_name,
1446 run_id,
1447 final_status,
1448 Some(err.to_string()),
1449 ctx,
1450 total_duration,
1451 run_labels,
1452 );
1453
1454 #[cfg(feature = "prometheus")]
1455 self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
1456
1457 return Err(err);
1458 }
1459 }
1460
1461 self.publish_run_status_changed(
1462 workflow_name,
1463 run_id,
1464 final_status,
1465 None,
1466 ctx,
1467 total_duration,
1468 run_labels,
1469 );
1470
1471 #[cfg(feature = "prometheus")]
1472 self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
1473
1474 Ok(WorkflowResult {
1475 run: final_run,
1476 steps: ctx.step_results().to_vec(),
1477 })
1478 }
1479
1480 #[cfg(feature = "prometheus")]
1482 fn emit_run_metrics(
1483 &self,
1484 workflow_name: &str,
1485 status: RunStatus,
1486 duration_ms: u64,
1487 ctx: &WorkflowContext,
1488 ) {
1489 let status_str = status.to_string();
1490 let wf = workflow_name.to_string();
1491
1492 counter!(RUNS_TOTAL, "workflow" => wf.clone(), "status" => status_str.clone()).increment(1);
1493 histogram!(RUN_DURATION_SECONDS, "workflow" => wf.clone(), "status" => status_str)
1494 .record(duration_ms as f64 / 1000.0);
1495 histogram!(RUN_COST_USD, "workflow" => wf.clone()).record(
1496 ctx.total_cost_usd()
1497 .to_string()
1498 .parse::<f64>()
1499 .unwrap_or(0.0),
1500 );
1501 gauge!(RUNS_ACTIVE, "workflow" => wf).decrement(1.0);
1502 }
1503
1504 fn on_run_budget_exceeded(&self, workflow_name: &str, run_id: Uuid, err: &EngineError) {
1510 let EngineError::RunBudgetExceeded {
1511 limit_usd,
1512 spent_usd,
1513 step_budget_usd,
1514 ..
1515 } = err
1516 else {
1517 return;
1518 };
1519
1520 #[cfg(feature = "prometheus")]
1521 counter!(
1522 RUN_BUDGET_EXCEEDED_TOTAL,
1523 "workflow" => workflow_name.to_string(),
1524 "scope" => "run",
1525 )
1526 .increment(1);
1527
1528 self.event_publisher
1529 .publish(Event::RunBudgetExceeded(RunBudgetExceededEvent {
1530 run_id,
1531 workflow_name: workflow_name.to_string(),
1532 limit_usd: *limit_usd,
1533 spent_usd: *spent_usd,
1534 step_budget_usd: *step_budget_usd,
1535 at: Utc::now(),
1536 }));
1537 }
1538
1539 #[allow(clippy::too_many_arguments)]
1544 fn publish_run_status_changed(
1545 &self,
1546 workflow_name: &str,
1547 run_id: Uuid,
1548 to: RunStatus,
1549 error: Option<String>,
1550 ctx: &WorkflowContext,
1551 duration_ms: u64,
1552 labels: HashMap<String, String>,
1553 ) {
1554 let now = Utc::now();
1555 let cost_usd = ctx.total_cost_usd();
1556 let wf = workflow_name.to_string();
1557
1558 self.event_publisher
1559 .publish(Event::RunStatusChanged(RunStatusChangedEvent {
1560 run_id,
1561 workflow_name: wf.clone(),
1562 from: RunStatus::Running,
1563 to,
1564 error: error.clone(),
1565 cost_usd,
1566 duration_ms,
1567 labels: labels.clone(),
1568 at: now,
1569 }));
1570
1571 if to == RunStatus::Failed {
1572 self.event_publisher
1573 .publish(Event::RunFailed(RunFailedEvent {
1574 run_id,
1575 workflow_name: wf,
1576 error,
1577 cost_usd,
1578 duration_ms,
1579 labels,
1580 at: now,
1581 }));
1582 }
1583 }
1584}
1585
1586impl fmt::Debug for Engine {
1587 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1588 f.debug_struct("Engine")
1589 .field("handlers", &self.handlers.keys().collect::<Vec<_>>())
1590 .finish_non_exhaustive()
1591 }
1592}
1593
1594#[cfg(test)]
1595mod tests {
1596 use super::*;
1597 use crate::config::ShellConfig;
1598 use crate::handler::{HandlerFuture, WorkflowHandler};
1599 use ironflow_core::providers::claude::ClaudeCodeProvider;
1600 use ironflow_core::providers::record_replay::RecordReplayProvider;
1601 use ironflow_store::memory::InMemoryStore;
1602 use ironflow_store::models::StepStatus;
1603 use serde_json::json;
1604
1605 struct EchoWorkflow;
1607
1608 impl WorkflowHandler for EchoWorkflow {
1609 fn name(&self) -> &str {
1610 "echo-workflow"
1611 }
1612
1613 fn describe(&self) -> WorkflowInfo {
1614 WorkflowInfo {
1615 description: "A simple workflow that echoes hello".to_string(),
1616 source_code: None,
1617 sub_workflows: Vec::new(),
1618 category: None,
1619 version: self.version().map(str::to_string),
1620 compatible_versions: Vec::new(),
1621 input_schema: None,
1622 default_labels: HashMap::new(),
1623 schedule: self.schedule().cloned(),
1624 default_max_cost_usd: self.default_max_cost_usd(),
1625 }
1626 }
1627
1628 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
1629 Box::pin(async move {
1630 ctx.shell("greet", ShellConfig::new("echo hello")).await?;
1631 Ok(())
1632 })
1633 }
1634 }
1635
1636 struct FailingWorkflow;
1638
1639 impl WorkflowHandler for FailingWorkflow {
1640 fn name(&self) -> &str {
1641 "failing-workflow"
1642 }
1643
1644 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
1645 Box::pin(async move {
1646 ctx.shell("fail", ShellConfig::new("exit 1")).await?;
1647 Ok(())
1648 })
1649 }
1650 }
1651
1652 fn create_test_engine() -> Engine {
1653 let store = Arc::new(InMemoryStore::new());
1654 let inner = ClaudeCodeProvider::new();
1655 let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
1656 inner,
1657 "/tmp/ironflow-fixtures",
1658 ));
1659 Engine::new(store, provider)
1660 }
1661
1662 #[test]
1663 fn engine_new_creates_instance() {
1664 let engine = create_test_engine();
1665 assert_eq!(engine.handler_names().len(), 0);
1666 }
1667
1668 #[test]
1669 fn engine_register_handler() {
1670 let mut engine = create_test_engine();
1671 let result = engine.register(EchoWorkflow);
1672 assert!(result.is_ok());
1673 assert_eq!(engine.handler_names().len(), 1);
1674 assert!(engine.handler_names().contains(&"echo-workflow"));
1675 }
1676
1677 #[test]
1678 fn engine_register_duplicate_returns_error() {
1679 let mut engine = create_test_engine();
1680 engine.register(EchoWorkflow).unwrap();
1681 let result = engine.register(EchoWorkflow);
1682 assert!(result.is_err());
1683 }
1684
1685 #[test]
1686 fn engine_get_handler_found() {
1687 let mut engine = create_test_engine();
1688 engine.register(EchoWorkflow).unwrap();
1689 let handler = engine.get_handler("echo-workflow");
1690 assert!(handler.is_some());
1691 }
1692
1693 #[test]
1694 fn engine_get_handler_not_found() {
1695 let engine = create_test_engine();
1696 let handler = engine.get_handler("nonexistent");
1697 assert!(handler.is_none());
1698 }
1699
1700 #[test]
1701 fn engine_handler_names_lists_all() {
1702 let mut engine = create_test_engine();
1703 engine.register(EchoWorkflow).unwrap();
1704 engine.register(FailingWorkflow).unwrap();
1705 let names = engine.handler_names();
1706 assert_eq!(names.len(), 2);
1707 assert!(names.contains(&"echo-workflow"));
1708 assert!(names.contains(&"failing-workflow"));
1709 }
1710
1711 #[test]
1712 fn engine_handler_info_returns_description() {
1713 let mut engine = create_test_engine();
1714 engine.register(EchoWorkflow).unwrap();
1715 let info = engine.handler_info("echo-workflow");
1716 assert!(info.is_some());
1717 let info = info.unwrap();
1718 assert_eq!(info.description, "A simple workflow that echoes hello");
1719 }
1720
1721 struct CategorizedWorkflow;
1722
1723 impl WorkflowHandler for CategorizedWorkflow {
1724 fn name(&self) -> &str {
1725 "categorized"
1726 }
1727 fn category(&self) -> Option<&str> {
1728 Some("data/etl")
1729 }
1730 fn execute<'a>(
1731 &'a self,
1732 _ctx: &'a mut WorkflowContext,
1733 ) -> crate::handler::HandlerFuture<'a> {
1734 Box::pin(async move { Ok(()) })
1735 }
1736 }
1737
1738 #[test]
1739 fn engine_default_describe_propagates_category() {
1740 let mut engine = create_test_engine();
1741 engine.register(CategorizedWorkflow).unwrap();
1742 let info = engine.handler_info("categorized").unwrap();
1743 assert_eq!(info.category.as_deref(), Some("data/etl"));
1744 }
1745
1746 #[test]
1747 fn engine_default_describe_without_category() {
1748 let mut engine = create_test_engine();
1749 engine.register(EchoWorkflow).unwrap();
1750 let info = engine.handler_info("echo-workflow").unwrap();
1751 assert!(info.category.is_none());
1752 }
1753
1754 struct ScheduledWorkflow {
1759 schedule: CronSchedule,
1760 }
1761
1762 impl ScheduledWorkflow {
1763 fn new() -> Self {
1764 Self {
1765 schedule: CronSchedule::new("0 0 * * * *").unwrap(),
1766 }
1767 }
1768 }
1769
1770 impl WorkflowHandler for ScheduledWorkflow {
1771 fn name(&self) -> &str {
1772 "scheduled"
1773 }
1774 fn schedule(&self) -> Option<&CronSchedule> {
1775 Some(&self.schedule)
1776 }
1777 fn execute<'a>(
1778 &'a self,
1779 _ctx: &'a mut WorkflowContext,
1780 ) -> crate::handler::HandlerFuture<'a> {
1781 Box::pin(async move { Ok(()) })
1782 }
1783 }
1784
1785 #[test]
1786 fn engine_default_describe_propagates_schedule() {
1787 let mut engine = create_test_engine();
1788 engine.register(ScheduledWorkflow::new()).unwrap();
1789 let info = engine.handler_info("scheduled").unwrap();
1790 assert_eq!(
1791 info.schedule.as_ref().map(|s| s.as_str()),
1792 Some("0 0 * * * *")
1793 );
1794 }
1795
1796 #[test]
1797 fn engine_default_describe_without_schedule() {
1798 let mut engine = create_test_engine();
1799 engine.register(EchoWorkflow).unwrap();
1800 let info = engine.handler_info("echo-workflow").unwrap();
1801 assert!(info.schedule.is_none());
1802 }
1803
1804 #[test]
1805 fn scheduled_handlers_returns_only_scheduled() {
1806 let mut engine = create_test_engine();
1807 engine.register(EchoWorkflow).unwrap();
1808 engine.register(ScheduledWorkflow::new()).unwrap();
1809 engine.register(FailingWorkflow).unwrap();
1810
1811 let scheduled = engine.scheduled_handlers();
1812 assert_eq!(scheduled.len(), 1);
1813 assert_eq!(scheduled[0].0, "scheduled");
1814 assert_eq!(scheduled[0].1.as_str(), "0 0 * * * *");
1815 }
1816
1817 #[test]
1818 fn scheduled_handlers_empty_when_none_scheduled() {
1819 let mut engine = create_test_engine();
1820 engine.register(EchoWorkflow).unwrap();
1821 engine.register(FailingWorkflow).unwrap();
1822
1823 let scheduled = engine.scheduled_handlers();
1824 assert!(scheduled.is_empty());
1825 }
1826
1827 struct BadCategoryWorkflow(&'static str);
1828
1829 impl WorkflowHandler for BadCategoryWorkflow {
1830 fn name(&self) -> &str {
1831 "bad-category"
1832 }
1833 fn category(&self) -> Option<&str> {
1834 Some(self.0)
1835 }
1836 fn execute<'a>(
1837 &'a self,
1838 _ctx: &'a mut WorkflowContext,
1839 ) -> crate::handler::HandlerFuture<'a> {
1840 Box::pin(async move { Ok(()) })
1841 }
1842 }
1843
1844 #[test]
1845 fn engine_register_rejects_empty_category() {
1846 let mut engine = create_test_engine();
1847 let err = engine.register(BadCategoryWorkflow("")).unwrap_err();
1848 match err {
1849 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty category")),
1850 other => panic!("expected InvalidWorkflow, got {other:?}"),
1851 }
1852 }
1853
1854 #[test]
1855 fn engine_register_rejects_leading_slash_category() {
1856 let mut engine = create_test_engine();
1857 let err = engine
1858 .register(BadCategoryWorkflow("/data/etl"))
1859 .unwrap_err();
1860 match err {
1861 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("leading '/'")),
1862 other => panic!("expected InvalidWorkflow, got {other:?}"),
1863 }
1864 }
1865
1866 #[test]
1867 fn engine_register_rejects_trailing_slash_category() {
1868 let mut engine = create_test_engine();
1869 let err = engine
1870 .register(BadCategoryWorkflow("data/etl/"))
1871 .unwrap_err();
1872 match err {
1873 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("trailing '/'")),
1874 other => panic!("expected InvalidWorkflow, got {other:?}"),
1875 }
1876 }
1877
1878 #[test]
1879 fn engine_register_rejects_double_slash_category() {
1880 let mut engine = create_test_engine();
1881 let err = engine
1882 .register(BadCategoryWorkflow("data//etl"))
1883 .unwrap_err();
1884 match err {
1885 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty segment")),
1886 other => panic!("expected InvalidWorkflow, got {other:?}"),
1887 }
1888 }
1889
1890 #[test]
1891 fn engine_register_rejects_whitespace_only_segment_category() {
1892 let mut engine = create_test_engine();
1893 let err = engine
1894 .register(BadCategoryWorkflow("data/ /etl"))
1895 .unwrap_err();
1896 match err {
1897 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("whitespace-only segment")),
1898 other => panic!("expected InvalidWorkflow, got {other:?}"),
1899 }
1900 }
1901
1902 #[test]
1903 fn engine_register_accepts_valid_nested_category() {
1904 let mut engine = create_test_engine();
1905 assert!(engine.register(CategorizedWorkflow).is_ok());
1906 }
1907
1908 #[tokio::test]
1909 async fn engine_unknown_workflow_returns_error() {
1910 let engine = create_test_engine();
1911 let result = engine
1912 .run_handler("unknown", TriggerKind::Manual, json!({}))
1913 .await;
1914 assert!(result.is_err());
1915 match result {
1916 Err(EngineError::InvalidWorkflow(msg)) => {
1917 assert!(msg.contains("no handler registered"));
1918 }
1919 _ => panic!("expected InvalidWorkflow error"),
1920 }
1921 }
1922
1923 #[tokio::test]
1924 async fn engine_enqueue_handler_creates_pending_run() {
1925 let mut engine = create_test_engine();
1926 engine.register(EchoWorkflow).unwrap();
1927
1928 let run = engine
1929 .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
1930 .await
1931 .unwrap();
1932 assert_eq!(run.status.state, RunStatus::Pending);
1933 assert_eq!(run.workflow_name, "echo-workflow");
1934 }
1935
1936 #[tokio::test]
1937 async fn enqueue_handler_leaves_the_run_unattributed() {
1938 let mut engine = create_test_engine();
1939 engine.register(EchoWorkflow).unwrap();
1940
1941 let run = engine
1942 .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
1943 .await
1944 .unwrap();
1945
1946 assert!(run.created_by.is_none());
1947 }
1948
1949 #[tokio::test]
1950 async fn enqueue_handler_with_options_records_the_author() {
1951 let mut engine = create_test_engine();
1952 engine.register(EchoWorkflow).unwrap();
1953 let actor = RunActor::User {
1954 user_id: Uuid::now_v7(),
1955 };
1956
1957 let run = engine
1958 .enqueue_handler_with_options(
1959 "echo-workflow",
1960 TriggerKind::Api,
1961 json!({}),
1962 EnqueueOptions {
1963 created_by: Some(actor.clone()),
1964 ..Default::default()
1965 },
1966 )
1967 .await
1968 .unwrap()
1969 .into_run();
1970
1971 assert_eq!(run.created_by, Some(actor));
1972 }
1973
1974 #[tokio::test]
1975 async fn enqueue_handler_with_options_accepts_no_author() {
1976 let mut engine = create_test_engine();
1977 engine.register(EchoWorkflow).unwrap();
1978
1979 let run = engine
1980 .enqueue_handler_with_options(
1981 "echo-workflow",
1982 TriggerKind::Cron {
1983 schedule: "0 * * * * *".to_string(),
1984 },
1985 json!({}),
1986 EnqueueOptions::default(),
1987 )
1988 .await
1989 .unwrap()
1990 .into_run();
1991
1992 assert!(run.created_by.is_none());
1993 }
1994
1995 #[tokio::test]
1996 async fn run_handler_leaves_the_run_unattributed() {
1997 let mut engine = create_test_engine();
1998 engine.register(EchoWorkflow).unwrap();
1999
2000 let run = engine
2001 .run_handler("echo-workflow", TriggerKind::Manual, json!({}))
2002 .await
2003 .unwrap()
2004 .run;
2005
2006 assert!(run.created_by.is_none());
2007 }
2008
2009 #[tokio::test]
2010 async fn engine_register_boxed() {
2011 let mut engine = create_test_engine();
2012 let handler: Box<dyn WorkflowHandler> = Box::new(EchoWorkflow);
2013 let result = engine.register_boxed(handler);
2014 assert!(result.is_ok());
2015 assert_eq!(engine.handler_names().len(), 1);
2016 }
2017
2018 #[tokio::test]
2019 async fn engine_store_and_provider_accessors() {
2020 let store = Arc::new(InMemoryStore::new());
2021 let inner = ClaudeCodeProvider::new();
2022 let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
2023 inner,
2024 "/tmp/ironflow-fixtures",
2025 ));
2026 let engine = Engine::new(store.clone(), provider.clone());
2027
2028 let _ = engine.store();
2030 let _ = engine.provider();
2031 }
2032
2033 use crate::operation::{Operation, OperationContext};
2038 use async_trait::async_trait;
2039 use ironflow_core::error::OperationError;
2040 use ironflow_store::models::StepKind;
2041
2042 struct FakeGitlabOp {
2043 project_id: u64,
2044 title: String,
2045 }
2046
2047 #[async_trait]
2048 impl Operation for FakeGitlabOp {
2049 fn kind(&self) -> &str {
2050 "gitlab"
2051 }
2052
2053 async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
2054 Ok(json!({
2055 "issue_id": 42,
2056 "project_id": self.project_id,
2057 "title": self.title,
2058 }))
2059 }
2060
2061 fn input(&self) -> Option<Value> {
2062 Some(json!({
2063 "project_id": self.project_id,
2064 "title": self.title,
2065 }))
2066 }
2067 }
2068
2069 struct FailingOp;
2070
2071 #[async_trait]
2072 impl Operation for FailingOp {
2073 fn kind(&self) -> &str {
2074 "broken-service"
2075 }
2076
2077 async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
2078 Err(OperationError::Http {
2079 status: None,
2080 message: "service unavailable".to_string(),
2081 })
2082 }
2083 }
2084
2085 struct OperationWorkflow;
2086
2087 impl WorkflowHandler for OperationWorkflow {
2088 fn name(&self) -> &str {
2089 "operation-workflow"
2090 }
2091
2092 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2093 Box::pin(async move {
2094 let op = FakeGitlabOp {
2095 project_id: 123,
2096 title: "Bug report".to_string(),
2097 };
2098 ctx.operation("create-issue", &op).await?;
2099 Ok(())
2100 })
2101 }
2102 }
2103
2104 struct FailingOperationWorkflow;
2105
2106 impl WorkflowHandler for FailingOperationWorkflow {
2107 fn name(&self) -> &str {
2108 "failing-operation-workflow"
2109 }
2110
2111 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2112 Box::pin(async move {
2113 ctx.operation("broken-call", &FailingOp).await?;
2114 Ok(())
2115 })
2116 }
2117 }
2118
2119 struct MixedWorkflow;
2120
2121 impl WorkflowHandler for MixedWorkflow {
2122 fn name(&self) -> &str {
2123 "mixed-workflow"
2124 }
2125
2126 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2127 Box::pin(async move {
2128 ctx.shell("build", ShellConfig::new("echo built")).await?;
2129 let op = FakeGitlabOp {
2130 project_id: 456,
2131 title: "Deploy done".to_string(),
2132 };
2133 let result = ctx.operation("notify-gitlab", &op).await?;
2134 assert_eq!(result.output["issue_id"], 42);
2135 Ok(())
2136 })
2137 }
2138 }
2139
2140 #[tokio::test]
2141 async fn operation_step_happy_path() {
2142 let mut engine = create_test_engine();
2143 engine.register(OperationWorkflow).unwrap();
2144
2145 let run = engine
2146 .run_handler("operation-workflow", TriggerKind::Manual, json!({}))
2147 .await
2148 .unwrap()
2149 .run;
2150
2151 assert_eq!(run.status.state, RunStatus::Completed);
2152
2153 let steps = engine.store().list_steps(run.id).await.unwrap();
2154
2155 assert_eq!(steps.len(), 1);
2156 assert_eq!(steps[0].name, "create-issue");
2157 assert_eq!(steps[0].kind, StepKind::Custom("gitlab".to_string()));
2158 assert_eq!(
2159 steps[0].status.state,
2160 ironflow_store::models::StepStatus::Completed
2161 );
2162
2163 let output = steps[0].output.as_ref().unwrap();
2164 assert_eq!(output["issue_id"], 42);
2165 assert_eq!(output["project_id"], 123);
2166
2167 let input = steps[0].input.as_ref().unwrap();
2168 assert_eq!(input["project_id"], 123);
2169 assert_eq!(input["title"], "Bug report");
2170 }
2171
2172 #[tokio::test]
2173 async fn operation_step_failure_marks_run_failed() {
2174 let mut engine = create_test_engine();
2175 engine.register(FailingOperationWorkflow).unwrap();
2176
2177 let result = engine
2178 .run_handler("failing-operation-workflow", TriggerKind::Manual, json!({}))
2179 .await;
2180
2181 assert!(result.is_err());
2182 }
2183
2184 #[tokio::test]
2185 async fn operation_mixed_with_shell_steps() {
2186 let mut engine = create_test_engine();
2187 engine.register(MixedWorkflow).unwrap();
2188
2189 let run = engine
2190 .run_handler("mixed-workflow", TriggerKind::Manual, json!({}))
2191 .await
2192 .unwrap()
2193 .run;
2194
2195 assert_eq!(run.status.state, RunStatus::Completed);
2196
2197 let steps = engine.store().list_steps(run.id).await.unwrap();
2198
2199 assert_eq!(steps.len(), 2);
2200 assert_eq!(steps[0].kind, StepKind::Shell);
2201 assert_eq!(steps[1].kind, StepKind::Custom("gitlab".to_string()));
2202 assert_eq!(steps[0].position, 0);
2203 assert_eq!(steps[1].position, 1);
2204 }
2205
2206 use crate::config::ApprovalConfig;
2211
2212 struct SingleApprovalWorkflow;
2213
2214 impl WorkflowHandler for SingleApprovalWorkflow {
2215 fn name(&self) -> &str {
2216 "single-approval"
2217 }
2218
2219 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2220 Box::pin(async move {
2221 ctx.shell("build", ShellConfig::new("echo built")).await?;
2222 ctx.approval("gate", ApprovalConfig::new("OK?")).await?;
2223 ctx.shell("deploy", ShellConfig::new("echo deployed"))
2224 .await?;
2225 Ok(())
2226 })
2227 }
2228 }
2229
2230 struct DoubleApprovalWorkflow;
2231
2232 impl WorkflowHandler for DoubleApprovalWorkflow {
2233 fn name(&self) -> &str {
2234 "double-approval"
2235 }
2236
2237 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2238 Box::pin(async move {
2239 ctx.shell("build", ShellConfig::new("echo built")).await?;
2240 ctx.approval("staging-gate", ApprovalConfig::new("Deploy staging?"))
2241 .await?;
2242 ctx.shell("deploy-staging", ShellConfig::new("echo staging"))
2243 .await?;
2244 ctx.approval("prod-gate", ApprovalConfig::new("Deploy prod?"))
2245 .await?;
2246 ctx.shell("deploy-prod", ShellConfig::new("echo prod"))
2247 .await?;
2248 Ok(())
2249 })
2250 }
2251 }
2252
2253 #[tokio::test]
2254 async fn approval_pauses_run() {
2255 let mut engine = create_test_engine();
2256 engine.register(SingleApprovalWorkflow).unwrap();
2257
2258 let run = engine
2259 .run_handler("single-approval", TriggerKind::Manual, json!({}))
2260 .await
2261 .unwrap()
2262 .run;
2263
2264 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
2265
2266 let steps = engine.store().list_steps(run.id).await.unwrap();
2267 assert_eq!(steps.len(), 2); assert_eq!(steps[0].kind, StepKind::Shell);
2269 assert_eq!(steps[0].status.state, StepStatus::Completed);
2270 assert_eq!(steps[1].kind, StepKind::Approval);
2271 assert_eq!(steps[1].status.state, StepStatus::AwaitingApproval);
2272 }
2273
2274 #[tokio::test]
2275 async fn approval_resume_completes_run() {
2276 let mut engine = create_test_engine();
2277 engine.register(SingleApprovalWorkflow).unwrap();
2278
2279 let run = engine
2281 .run_handler("single-approval", TriggerKind::Manual, json!({}))
2282 .await
2283 .unwrap()
2284 .run;
2285 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
2286
2287 engine
2289 .store()
2290 .update_run_status(run.id, RunStatus::Running)
2291 .await
2292 .unwrap();
2293
2294 let resumed = engine.resume_run(run.id).await.unwrap().run;
2296 assert_eq!(resumed.status.state, RunStatus::Completed);
2297
2298 let steps = engine.store().list_steps(run.id).await.unwrap();
2299 assert_eq!(steps.len(), 3); assert_eq!(steps[0].name, "build");
2301 assert_eq!(steps[0].status.state, StepStatus::Completed);
2302 assert_eq!(steps[1].name, "gate");
2303 assert_eq!(steps[1].kind, StepKind::Approval);
2304 assert_eq!(steps[1].status.state, StepStatus::Completed);
2305 assert_eq!(steps[2].name, "deploy");
2306 assert_eq!(steps[2].status.state, StepStatus::Completed);
2307 }
2308
2309 #[tokio::test]
2310 async fn double_approval_two_resumes() {
2311 let mut engine = create_test_engine();
2312 engine.register(DoubleApprovalWorkflow).unwrap();
2313
2314 let run = engine
2316 .run_handler("double-approval", TriggerKind::Manual, json!({}))
2317 .await
2318 .unwrap()
2319 .run;
2320 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
2321
2322 let steps = engine.store().list_steps(run.id).await.unwrap();
2323 assert_eq!(steps.len(), 2); engine
2327 .store()
2328 .update_run_status(run.id, RunStatus::Running)
2329 .await
2330 .unwrap();
2331
2332 let resumed = engine.resume_run(run.id).await.unwrap().run;
2333 assert_eq!(resumed.status.state, RunStatus::AwaitingApproval);
2334
2335 let steps = engine.store().list_steps(run.id).await.unwrap();
2336 assert_eq!(steps.len(), 4); engine
2340 .store()
2341 .update_run_status(run.id, RunStatus::Running)
2342 .await
2343 .unwrap();
2344
2345 let final_run = engine.resume_run(run.id).await.unwrap().run;
2346 assert_eq!(final_run.status.state, RunStatus::Completed);
2347
2348 let steps = engine.store().list_steps(run.id).await.unwrap();
2349 assert_eq!(steps.len(), 5);
2350 assert_eq!(steps[0].name, "build");
2351 assert_eq!(steps[1].name, "staging-gate");
2352 assert_eq!(steps[2].name, "deploy-staging");
2353 assert_eq!(steps[3].name, "prod-gate");
2354 assert_eq!(steps[4].name, "deploy-prod");
2355
2356 for step in &steps {
2357 assert_eq!(step.status.state, StepStatus::Completed);
2358 }
2359 }
2360
2361 use ironflow_store::models::{NewStep, StepUpdate, step_trace_id};
2366
2367 async fn create_step_with_status(
2368 store: &Arc<dyn Store>,
2369 run_id: Uuid,
2370 name: &str,
2371 position: u32,
2372 status: StepStatus,
2373 ) -> ironflow_store::models::Step {
2374 let step = store
2375 .create_step(NewStep {
2376 run_id,
2377 trace_id: step_trace_id(run_id, name, position),
2378 name: name.to_string(),
2379 kind: StepKind::Shell,
2380 position,
2381 input: None,
2382 is_error_handler: false,
2383 })
2384 .await
2385 .unwrap();
2386
2387 match status {
2388 StepStatus::Pending => {}
2389 StepStatus::Running => {
2390 store
2391 .update_step(
2392 step.id,
2393 StepUpdate {
2394 status: Some(StepStatus::Running),
2395 ..StepUpdate::default()
2396 },
2397 )
2398 .await
2399 .unwrap();
2400 }
2401 StepStatus::Completed => {
2402 store
2403 .update_step(
2404 step.id,
2405 StepUpdate {
2406 status: Some(StepStatus::Running),
2407 ..StepUpdate::default()
2408 },
2409 )
2410 .await
2411 .unwrap();
2412 store
2413 .update_step(
2414 step.id,
2415 StepUpdate {
2416 status: Some(StepStatus::Completed),
2417 ..StepUpdate::default()
2418 },
2419 )
2420 .await
2421 .unwrap();
2422 }
2423 StepStatus::AwaitingApproval => {
2424 store
2425 .update_step(
2426 step.id,
2427 StepUpdate {
2428 status: Some(StepStatus::Running),
2429 ..StepUpdate::default()
2430 },
2431 )
2432 .await
2433 .unwrap();
2434 store
2435 .update_step(
2436 step.id,
2437 StepUpdate {
2438 status: Some(StepStatus::AwaitingApproval),
2439 ..StepUpdate::default()
2440 },
2441 )
2442 .await
2443 .unwrap();
2444 }
2445 _ => panic!("unsupported status for test helper: {status}"),
2446 }
2447
2448 store.get_step(step.id).await.unwrap().unwrap()
2449 }
2450
2451 #[tokio::test]
2452 async fn fail_orphaned_steps_marks_running_as_failed() {
2453 let engine = create_test_engine();
2454 let run = engine
2455 .store()
2456 .create_run(NewRun {
2457 created_by: None,
2458 workflow_name: "test".to_string(),
2459 trigger: TriggerKind::Manual,
2460 payload: json!({}),
2461 max_retries: 0,
2462 handler_version: None,
2463 labels: HashMap::new(),
2464 scheduled_at: None,
2465 idempotency_key: None,
2466 max_cost_usd: None,
2467 })
2468 .await
2469 .unwrap()
2470 .into_run();
2471
2472 let step = create_step_with_status(
2473 engine.store(),
2474 run.id,
2475 "running-step",
2476 0,
2477 StepStatus::Running,
2478 )
2479 .await;
2480
2481 engine
2482 .fail_orphaned_steps(run.id, "parent run timed out")
2483 .await
2484 .unwrap();
2485
2486 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
2487 assert_eq!(updated.status.state, StepStatus::Failed);
2488 assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
2489 assert!(updated.completed_at.is_some());
2490 }
2491
2492 #[tokio::test]
2493 async fn fail_orphaned_steps_marks_pending_as_skipped() {
2494 let engine = create_test_engine();
2495 let run = engine
2496 .store()
2497 .create_run(NewRun {
2498 created_by: None,
2499 workflow_name: "test".to_string(),
2500 trigger: TriggerKind::Manual,
2501 payload: json!({}),
2502 max_retries: 0,
2503 handler_version: None,
2504 labels: HashMap::new(),
2505 scheduled_at: None,
2506 idempotency_key: None,
2507 max_cost_usd: None,
2508 })
2509 .await
2510 .unwrap()
2511 .into_run();
2512
2513 let step = create_step_with_status(
2514 engine.store(),
2515 run.id,
2516 "pending-step",
2517 0,
2518 StepStatus::Pending,
2519 )
2520 .await;
2521
2522 engine
2523 .fail_orphaned_steps(run.id, "parent run timed out")
2524 .await
2525 .unwrap();
2526
2527 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
2528 assert_eq!(updated.status.state, StepStatus::Skipped);
2529 assert!(updated.error.is_none());
2530 assert!(updated.completed_at.is_some());
2531 }
2532
2533 #[tokio::test]
2534 async fn fail_orphaned_steps_marks_awaiting_approval_as_failed() {
2535 let engine = create_test_engine();
2536 let run = engine
2537 .store()
2538 .create_run(NewRun {
2539 created_by: None,
2540 workflow_name: "test".to_string(),
2541 trigger: TriggerKind::Manual,
2542 payload: json!({}),
2543 max_retries: 0,
2544 handler_version: None,
2545 labels: HashMap::new(),
2546 scheduled_at: None,
2547 idempotency_key: None,
2548 max_cost_usd: None,
2549 })
2550 .await
2551 .unwrap()
2552 .into_run();
2553
2554 let step = create_step_with_status(
2555 engine.store(),
2556 run.id,
2557 "approval-step",
2558 0,
2559 StepStatus::AwaitingApproval,
2560 )
2561 .await;
2562
2563 engine
2564 .fail_orphaned_steps(run.id, "parent run timed out")
2565 .await
2566 .unwrap();
2567
2568 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
2569 assert_eq!(updated.status.state, StepStatus::Failed);
2570 assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
2571 assert!(updated.completed_at.is_some());
2572 }
2573
2574 #[tokio::test]
2575 async fn fail_orphaned_steps_skips_terminal_steps() {
2576 let engine = create_test_engine();
2577 let run = engine
2578 .store()
2579 .create_run(NewRun {
2580 created_by: None,
2581 workflow_name: "test".to_string(),
2582 trigger: TriggerKind::Manual,
2583 payload: json!({}),
2584 max_retries: 0,
2585 handler_version: None,
2586 labels: HashMap::new(),
2587 scheduled_at: None,
2588 idempotency_key: None,
2589 max_cost_usd: None,
2590 })
2591 .await
2592 .unwrap()
2593 .into_run();
2594
2595 let completed_step =
2596 create_step_with_status(engine.store(), run.id, "done", 0, StepStatus::Completed).await;
2597 let running_step =
2598 create_step_with_status(engine.store(), run.id, "in-flight", 1, StepStatus::Running)
2599 .await;
2600
2601 engine
2602 .fail_orphaned_steps(run.id, "parent run timed out")
2603 .await
2604 .unwrap();
2605
2606 let completed = engine
2607 .store()
2608 .get_step(completed_step.id)
2609 .await
2610 .unwrap()
2611 .unwrap();
2612 assert_eq!(completed.status.state, StepStatus::Completed);
2613
2614 let failed = engine
2615 .store()
2616 .get_step(running_step.id)
2617 .await
2618 .unwrap()
2619 .unwrap();
2620 assert_eq!(failed.status.state, StepStatus::Failed);
2621 }
2622
2623 #[tokio::test]
2624 async fn fail_orphaned_steps_mixed_states() {
2625 let engine = create_test_engine();
2626 let run = engine
2627 .store()
2628 .create_run(NewRun {
2629 created_by: None,
2630 workflow_name: "test".to_string(),
2631 trigger: TriggerKind::Manual,
2632 payload: json!({}),
2633 max_retries: 0,
2634 handler_version: None,
2635 labels: HashMap::new(),
2636 scheduled_at: None,
2637 idempotency_key: None,
2638 max_cost_usd: None,
2639 })
2640 .await
2641 .unwrap()
2642 .into_run();
2643
2644 let s_completed =
2645 create_step_with_status(engine.store(), run.id, "step-1", 0, StepStatus::Completed)
2646 .await;
2647 let s_running =
2648 create_step_with_status(engine.store(), run.id, "step-2", 1, StepStatus::Running).await;
2649 let s_pending =
2650 create_step_with_status(engine.store(), run.id, "step-3", 2, StepStatus::Pending).await;
2651
2652 engine.fail_orphaned_steps(run.id, "timeout").await.unwrap();
2653
2654 let r_completed = engine
2655 .store()
2656 .get_step(s_completed.id)
2657 .await
2658 .unwrap()
2659 .unwrap();
2660 assert_eq!(r_completed.status.state, StepStatus::Completed);
2661
2662 let r_running = engine
2663 .store()
2664 .get_step(s_running.id)
2665 .await
2666 .unwrap()
2667 .unwrap();
2668 assert_eq!(r_running.status.state, StepStatus::Failed);
2669 assert_eq!(r_running.error.as_deref(), Some("timeout"));
2670
2671 let r_pending = engine
2672 .store()
2673 .get_step(s_pending.id)
2674 .await
2675 .unwrap()
2676 .unwrap();
2677 assert_eq!(r_pending.status.state, StepStatus::Skipped);
2678 assert!(r_pending.error.is_none());
2679 }
2680
2681 #[tokio::test]
2682 async fn fail_orphaned_steps_no_steps_is_noop() {
2683 let engine = create_test_engine();
2684 let run = engine
2685 .store()
2686 .create_run(NewRun {
2687 created_by: None,
2688 workflow_name: "test".to_string(),
2689 trigger: TriggerKind::Manual,
2690 payload: json!({}),
2691 max_retries: 0,
2692 handler_version: None,
2693 labels: HashMap::new(),
2694 scheduled_at: None,
2695 idempotency_key: None,
2696 max_cost_usd: None,
2697 })
2698 .await
2699 .unwrap()
2700 .into_run();
2701
2702 let result = engine.fail_orphaned_steps(run.id, "timeout").await;
2703 assert!(result.is_ok());
2704 }
2705
2706 #[tokio::test]
2707 async fn fail_orphaned_steps_preserves_existing_error() {
2708 let engine = create_test_engine();
2709 let run = engine
2710 .store()
2711 .create_run(NewRun {
2712 created_by: None,
2713 workflow_name: "test".to_string(),
2714 trigger: TriggerKind::Manual,
2715 payload: json!({}),
2716 max_retries: 0,
2717 handler_version: None,
2718 labels: HashMap::new(),
2719 scheduled_at: None,
2720 idempotency_key: None,
2721 max_cost_usd: None,
2722 })
2723 .await
2724 .unwrap()
2725 .into_run();
2726
2727 let step_with_error = create_step_with_status(
2728 engine.store(),
2729 run.id,
2730 "already-errored",
2731 0,
2732 StepStatus::Running,
2733 )
2734 .await;
2735
2736 engine
2737 .store()
2738 .update_step(
2739 step_with_error.id,
2740 StepUpdate {
2741 error: Some("real error from provider".to_string()),
2742 ..StepUpdate::default()
2743 },
2744 )
2745 .await
2746 .unwrap();
2747
2748 let step_no_error = create_step_with_status(
2749 engine.store(),
2750 run.id,
2751 "no-error-yet",
2752 1,
2753 StepStatus::Running,
2754 )
2755 .await;
2756
2757 engine
2758 .fail_orphaned_steps(run.id, "parent run failed")
2759 .await
2760 .unwrap();
2761
2762 let updated_with = engine
2763 .store()
2764 .get_step(step_with_error.id)
2765 .await
2766 .unwrap()
2767 .unwrap();
2768 assert_eq!(updated_with.status.state, StepStatus::Failed);
2769 assert_eq!(
2770 updated_with.error.as_deref(),
2771 Some("real error from provider"),
2772 );
2773
2774 let updated_without = engine
2775 .store()
2776 .get_step(step_no_error.id)
2777 .await
2778 .unwrap()
2779 .unwrap();
2780 assert_eq!(updated_without.status.state, StepStatus::Failed);
2781 assert_eq!(updated_without.error.as_deref(), Some("parent run failed"),);
2782 }
2783}