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::{StepInterceptor, 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 step_interceptor: Option<Arc<dyn StepInterceptor>>,
173}
174
175fn validate_category(handler_name: &str, category: &str) -> Result<(), EngineError> {
185 let reject = |reason: &str| {
186 Err(EngineError::InvalidWorkflow(format!(
187 "handler '{handler_name}' has invalid category '{category}': {reason}"
188 )))
189 };
190
191 if category.is_empty() {
192 return reject("empty category");
193 }
194 if category.starts_with('/') {
195 return reject("leading '/'");
196 }
197 if category.ends_with('/') {
198 return reject("trailing '/'");
199 }
200 for segment in category.split('/') {
201 if segment.is_empty() {
202 return reject("empty segment (double '/')");
203 }
204 if segment.trim().is_empty() {
205 return reject("whitespace-only segment");
206 }
207 }
208 Ok(())
209}
210
211impl Engine {
212 pub fn new(store: Arc<dyn Store>, provider: Arc<dyn AgentProvider>) -> Self {
228 Self {
229 store,
230 provider,
231 handlers: HashMap::new(),
232 event_publisher: EventPublisher::new(),
233 log_sender: None,
234 budget: BudgetConfig::new(),
235 artifact_sink: None,
236 guard_config: None,
237 event_bus: None,
238 decision_provider: None,
239 step_interceptor: None,
240 }
241 }
242
243 pub fn with_decision_provider(mut self, provider: Arc<dyn DecisionProvider>) -> Self {
264 self.decision_provider = Some(provider);
265 self
266 }
267
268 pub fn with_step_interceptor(mut self, interceptor: Arc<dyn StepInterceptor>) -> Self {
292 self.step_interceptor = Some(interceptor);
293 self
294 }
295
296 pub fn step_interceptor(&self) -> Option<&Arc<dyn StepInterceptor>> {
298 self.step_interceptor.as_ref()
299 }
300
301 pub fn with_budget_config(mut self, budget: BudgetConfig) -> Self {
322 self.budget = budget;
323 self
324 }
325
326 pub fn budget_config(&self) -> &BudgetConfig {
328 &self.budget
329 }
330
331 pub fn with_guard_config(mut self, config: WorkflowGuardConfig) -> Self {
353 self.guard_config = Some(config);
354 self
355 }
356
357 pub fn guard_config(&self) -> Option<&WorkflowGuardConfig> {
359 self.guard_config.as_ref()
360 }
361
362 pub fn set_log_sender(&mut self, sender: LogSender) {
368 self.log_sender = Some(sender);
369 }
370
371 pub fn set_artifact_sink(&mut self, sink: Arc<dyn ArtifactSink>) {
389 self.artifact_sink = Some(sink);
390 }
391
392 pub fn artifact_sink(&self) -> Option<&Arc<dyn ArtifactSink>> {
394 self.artifact_sink.as_ref()
395 }
396
397 pub fn set_event_bus(&mut self, bus: WorkflowEventBus) {
414 self.event_bus = Some(bus);
415 }
416
417 pub fn event_bus(&self) -> Option<&WorkflowEventBus> {
419 self.event_bus.as_ref()
420 }
421
422 pub fn store(&self) -> &Arc<dyn Store> {
424 &self.store
425 }
426
427 pub fn provider(&self) -> &Arc<dyn AgentProvider> {
429 &self.provider
430 }
431
432 fn build_context(&self, run: &Run) -> WorkflowContext {
441 let handlers = self.handlers.clone();
442 let resolver: crate::context::HandlerResolver =
443 Arc::new(move |name: &str| handlers.get(name).cloned());
444 let mut ctx = WorkflowContext::with_handler_resolver(
445 run.id,
446 run.workflow_name.clone(),
447 self.store.clone(),
448 self.provider.clone(),
449 resolver,
450 );
451 ctx.carry_over_run_totals(run.retry_count + 1, run.cost_usd, run.duration_ms);
452 ctx.set_max_cost_usd(run.max_cost_usd);
453 if let Some(ref sender) = self.log_sender {
454 ctx.set_log_sender(sender.clone());
455 }
456 if let Some(ref sink) = self.artifact_sink {
457 ctx.set_artifact_sink(sink.clone());
458 }
459 if let Some(ref bus) = self.event_bus {
460 ctx.set_event_bus(bus.clone());
461 }
462 if let Some(ref provider) = self.decision_provider {
463 ctx.set_decision_provider(provider.clone());
464 }
465 if let Some(ref interceptor) = self.step_interceptor {
466 ctx.set_step_interceptor(interceptor.clone());
467 }
468 ctx
469 }
470
471 fn build_context_with_guard(
477 &self,
478 run: &Run,
479 handler: &dyn WorkflowHandler,
480 ) -> WorkflowContext {
481 let mut ctx = self.build_context(run);
482 let guard_config = handler.guard_config().or_else(|| self.guard_config.clone());
483 if let Some(config) = guard_config {
484 ctx.set_guard(config, new_shared_guard_state());
485 }
486 ctx
487 }
488
489 async fn check_monthly_quota(&self, workflow_name: &str) -> Result<(), EngineError> {
500 let Some(limit) = self.budget.monthly_cost_limit_usd else {
501 return Ok(());
502 };
503
504 let stats = self
505 .store
506 .get_stats(RunFilter {
507 created_after: Some(month_start(Utc::now())),
508 ..RunFilter::default()
509 })
510 .await?;
511
512 if stats.total_cost_usd < limit {
513 return Ok(());
514 }
515
516 warn!(
517 workflow = %workflow_name,
518 limit_usd = %limit,
519 spent_usd = %stats.total_cost_usd,
520 "monthly cost quota exhausted, refusing new run"
521 );
522
523 #[cfg(feature = "prometheus")]
524 counter!(
525 RUN_BUDGET_EXCEEDED_TOTAL,
526 "workflow" => workflow_name.to_string(),
527 "scope" => "monthly",
528 )
529 .increment(1);
530
531 Err(EngineError::MonthlyBudgetExceeded {
532 limit_usd: limit,
533 spent_usd: stats.total_cost_usd,
534 })
535 }
536
537 pub fn register(&mut self, handler: impl WorkflowHandler + 'static) -> Result<(), EngineError> {
581 let name = handler.name().to_string();
582 if self.handlers.contains_key(&name) {
583 return Err(EngineError::InvalidWorkflow(format!(
584 "handler '{}' already registered",
585 name
586 )));
587 }
588 if let Some(category) = handler.category() {
589 validate_category(&name, category)?;
590 }
591 self.handlers.insert(name, Arc::new(handler));
592 Ok(())
593 }
594
595 pub fn register_boxed(&mut self, handler: Box<dyn WorkflowHandler>) -> Result<(), EngineError> {
602 let name = handler.name().to_string();
603 if self.handlers.contains_key(&name) {
604 return Err(EngineError::InvalidWorkflow(format!(
605 "handler '{}' already registered",
606 name
607 )));
608 }
609 if let Some(category) = handler.category() {
610 validate_category(&name, category)?;
611 }
612 self.handlers.insert(name, Arc::from(handler));
613 Ok(())
614 }
615
616 pub fn get_handler(&self, name: &str) -> Option<&Arc<dyn WorkflowHandler>> {
618 self.handlers.get(name)
619 }
620
621 pub fn handler_names(&self) -> Vec<&str> {
623 self.handlers.keys().map(|s| s.as_str()).collect()
624 }
625
626 pub fn handler_info(&self, name: &str) -> Option<WorkflowInfo> {
628 self.handlers.get(name).map(|h| h.describe())
629 }
630
631 pub fn scheduled_handlers(&self) -> Vec<(&str, &CronSchedule)> {
655 self.handlers
656 .iter()
657 .filter_map(|(name, handler)| handler.schedule().map(|sched| (name.as_str(), sched)))
658 .collect()
659 }
660
661 pub fn subscribe(
686 &mut self,
687 subscriber: impl EventSubscriber + 'static,
688 event_types: &[&'static str],
689 ) {
690 self.event_publisher.subscribe(subscriber, event_types);
691 }
692
693 pub fn event_publisher(&self) -> &EventPublisher {
698 &self.event_publisher
699 }
700
701 #[tracing::instrument(name = "engine.run_handler", skip_all, fields(workflow = %handler_name))]
731 pub async fn run_handler(
732 &self,
733 handler_name: &str,
734 trigger: TriggerKind,
735 payload: Value,
736 ) -> Result<WorkflowResult, EngineError> {
737 let handler = self
738 .handlers
739 .get(handler_name)
740 .ok_or_else(|| {
741 EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
742 })?
743 .clone();
744
745 self.check_monthly_quota(handler_name).await?;
746
747 let handler_version = handler.version().map(str::to_string);
748 let max_cost_usd = self
749 .budget
750 .resolve_run_cap(None, handler.default_max_cost_usd());
751 let run = self
752 .store
753 .create_run(NewRun {
754 created_by: None,
755 workflow_name: handler_name.to_string(),
756 trigger,
757 payload,
758 max_retries: 0,
759 handler_version,
760 labels: handler.default_labels(),
761 scheduled_at: None,
762 idempotency_key: None,
763 max_cost_usd,
764 })
765 .await?
766 .into_run();
767
768 let run_id = run.id;
769 info!(run_id = %run_id, handler_version = run.handler_version.as_deref().unwrap_or(""), "run created");
770
771 self.store
772 .update_run_status(run_id, RunStatus::Running)
773 .await?;
774
775 #[cfg(feature = "prometheus")]
776 gauge!(RUNS_ACTIVE, "workflow" => handler_name.to_string()).increment(1.0);
777
778 let run_start = Instant::now();
779 let mut ctx = self.build_context_with_guard(&run, handler.as_ref());
780
781 let result = handler.execute(&mut ctx).await;
782 self.finalize_run(run_id, handler_name, result, &ctx, run_start, run.labels)
783 .await
784 }
785
786 #[tracing::instrument(name = "engine.plan_handler", skip_all, fields(workflow = %handler_name))]
828 pub async fn plan_handler(
829 &self,
830 handler_name: &str,
831 payload: Value,
832 options: PlanOptions,
833 ) -> Result<ExecutionPlan, EngineError> {
834 if options.max_depth == 0 {
835 return Err(EngineError::InvalidWorkflow(
836 "max_depth must be at least 1".to_string(),
837 ));
838 }
839
840 let handler = self
841 .handlers
842 .get(handler_name)
843 .ok_or_else(|| {
844 EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
845 })?
846 .clone();
847
848 let estimates = if options.estimate_durations {
849 estimate_durations(&self.store, handler_name, options.sample_runs).await?
850 } else {
851 HashMap::new()
852 };
853
854 let shared: SharedPlanRecorder = Arc::new(Mutex::new(PlanRecorder::new(
855 handler_name.to_string(),
856 payload,
857 options.max_depth,
858 estimates,
859 )));
860
861 let handlers = self.handlers.clone();
864 let resolver: crate::context::HandlerResolver =
865 Arc::new(move |name: &str| handlers.get(name).cloned());
866 let mut ctx = WorkflowContext::with_handler_resolver(
867 Uuid::now_v7(),
868 handler_name.to_string(),
869 self.store.clone(),
870 self.provider.clone(),
871 resolver,
872 );
873 ctx.set_plan(shared.clone());
874
875 if let Err(err) = handler.execute(&mut ctx).await {
876 lock_plan(&shared).fail(err.to_string());
877 }
878 drop(ctx);
879
880 let plan = match Arc::try_unwrap(shared) {
881 Ok(mutex) => mutex
882 .into_inner()
883 .unwrap_or_else(|poisoned| poisoned.into_inner())
884 .into_plan(),
885 Err(shared) => lock_plan(&shared).snapshot(),
886 };
887
888 info!(
889 workflow = %handler_name,
890 steps = plan.steps.len(),
891 truncated = plan.truncated,
892 "execution plan built"
893 );
894
895 Ok(plan)
896 }
897
898 #[tracing::instrument(name = "engine.enqueue_handler", skip_all, fields(workflow = %handler_name))]
909 pub async fn enqueue_handler(
910 &self,
911 handler_name: &str,
912 trigger: TriggerKind,
913 payload: Value,
914 max_retries: u32,
915 ) -> Result<Run, EngineError> {
916 self.enqueue_handler_with_options(
917 handler_name,
918 trigger,
919 payload,
920 EnqueueOptions {
921 max_retries,
922 ..Default::default()
923 },
924 )
925 .await
926 .map(RunCreation::into_run)
927 }
928
929 #[tracing::instrument(name = "engine.enqueue_handler_with_options", skip_all, fields(workflow = %handler_name))]
974 pub async fn enqueue_handler_with_options(
975 &self,
976 handler_name: &str,
977 trigger: TriggerKind,
978 payload: Value,
979 options: EnqueueOptions,
980 ) -> Result<RunCreation, EngineError> {
981 let EnqueueOptions {
982 max_retries,
983 labels,
984 scheduled_at,
985 max_cost_usd,
986 created_by,
987 idempotency_key,
988 } = options;
989
990 let handler = self.handlers.get(handler_name).ok_or_else(|| {
991 EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
992 })?;
993
994 self.check_monthly_quota(handler_name).await?;
995
996 let handler_version = handler.version().map(str::to_string);
997 let mut merged_labels = handler.default_labels();
998 merged_labels.extend(labels);
999 let resolved_cap = self
1000 .budget
1001 .resolve_run_cap(max_cost_usd, handler.default_max_cost_usd());
1002
1003 let creation = self
1004 .store
1005 .create_run(NewRun {
1006 workflow_name: handler_name.to_string(),
1007 trigger,
1008 payload,
1009 max_retries,
1010 handler_version,
1011 labels: merged_labels,
1012 scheduled_at,
1013 created_by,
1014 idempotency_key,
1015 max_cost_usd: resolved_cap,
1016 })
1017 .await?;
1018
1019 match &creation {
1020 RunCreation::Created(run) => info!(
1021 run_id = %run.id,
1022 workflow = %handler_name,
1023 max_cost_usd = ?resolved_cap,
1024 "handler run enqueued"
1025 ),
1026 RunCreation::Existing(run) => info!(
1027 run_id = %run.id,
1028 workflow = %handler_name,
1029 "idempotent replay, nothing enqueued"
1030 ),
1031 }
1032
1033 Ok(creation)
1034 }
1035
1036 #[tracing::instrument(name = "engine.execute_handler_run", skip_all, fields(run_id = %run_id))]
1045 pub async fn execute_handler_run(&self, run_id: Uuid) -> Result<WorkflowResult, EngineError> {
1046 let run = self
1047 .store
1048 .get_run(run_id)
1049 .await?
1050 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
1051
1052 let handler = self
1053 .handlers
1054 .get(&run.workflow_name)
1055 .ok_or_else(|| {
1056 EngineError::InvalidWorkflow(format!(
1057 "no handler registered: {}",
1058 run.workflow_name
1059 ))
1060 })?
1061 .clone();
1062
1063 #[cfg(feature = "prometheus")]
1064 gauge!(RUNS_ACTIVE, "workflow" => run.workflow_name.clone()).increment(1.0);
1065
1066 let run_start = Instant::now();
1067 let mut ctx = self.build_context_with_guard(&run, handler.as_ref());
1068
1069 if run.retry_count > 0 {
1072 ctx.load_replay_steps().await?;
1073 }
1074
1075 let result = handler.execute(&mut ctx).await;
1076 self.finalize_run(
1077 run_id,
1078 &run.workflow_name,
1079 result,
1080 &ctx,
1081 run_start,
1082 run.labels,
1083 )
1084 .await
1085 }
1086
1087 #[tracing::instrument(name = "engine.execute_run", skip_all, fields(run_id = %run_id))]
1095 pub async fn execute_run(&self, run_id: Uuid) -> Result<WorkflowResult, EngineError> {
1096 self.execute_handler_run(run_id).await
1097 }
1098
1099 #[tracing::instrument(name = "engine.resume_run", skip_all, fields(run_id = %run_id))]
1113 pub async fn resume_run(&self, run_id: Uuid) -> Result<WorkflowResult, EngineError> {
1114 let run = self
1115 .store
1116 .get_run(run_id)
1117 .await?
1118 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
1119
1120 let handler = self
1121 .handlers
1122 .get(&run.workflow_name)
1123 .ok_or_else(|| {
1124 EngineError::InvalidWorkflow(format!(
1125 "no handler registered: {}",
1126 run.workflow_name
1127 ))
1128 })?
1129 .clone();
1130
1131 info!(run_id = %run_id, workflow = %run.workflow_name, "resuming run after approval");
1132
1133 let run_start = Instant::now();
1134 let mut ctx = self.build_context_with_guard(&run, handler.as_ref());
1135 ctx.load_replay_steps().await?;
1136
1137 let result = handler.execute(&mut ctx).await;
1138 self.finalize_run(
1139 run_id,
1140 &run.workflow_name,
1141 result,
1142 &ctx,
1143 run_start,
1144 run.labels,
1145 )
1146 .await
1147 }
1148
1149 pub async fn fail_or_schedule_retry(
1194 &self,
1195 run_id: Uuid,
1196 error: &str,
1197 retryable: bool,
1198 cost_usd: Option<Decimal>,
1199 duration_ms: Option<u64>,
1200 ) -> Result<RunStatus, EngineError> {
1201 let run = self
1202 .store
1203 .get_run(run_id)
1204 .await?
1205 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
1206
1207 let has_attempts_left = run.retry_count < run.max_retries;
1208 let update = if retryable && has_attempts_left {
1209 let backoff = backoff_for_retry(run.retry_count);
1210 let scheduled_at = Utc::now() + TimeDelta::milliseconds(backoff.as_millis() as i64);
1211
1212 info!(
1213 run_id = %run_id,
1214 workflow = %run.workflow_name,
1215 attempt = run.retry_count + 1,
1216 max_retries = run.max_retries,
1217 backoff_secs = backoff.as_secs(),
1218 scheduled_at = %scheduled_at,
1219 "run failed, scheduling retry"
1220 );
1221
1222 RunUpdate {
1223 status: Some(RunStatus::Retrying),
1224 error: Some(error.to_string()),
1225 increment_retry: true,
1226 cost_usd,
1227 duration_ms,
1228 scheduled_at: Some(scheduled_at),
1229 ..RunUpdate::default()
1230 }
1231 } else {
1232 RunUpdate {
1233 status: Some(RunStatus::Failed),
1234 error: Some(error.to_string()),
1235 cost_usd,
1236 duration_ms,
1237 completed_at: Some(Utc::now()),
1238 ..RunUpdate::default()
1239 }
1240 };
1241
1242 let status = update.status.unwrap_or(RunStatus::Failed);
1243 self.store.update_run(run_id, update).await?;
1244 self.fail_orphaned_steps(run_id, error).await?;
1245
1246 Ok(status)
1247 }
1248
1249 pub async fn fail_orphaned_steps(
1263 &self,
1264 run_id: Uuid,
1265 error_message: &str,
1266 ) -> Result<(), EngineError> {
1267 let steps = self.store.list_steps(run_id).await?;
1268 let now = Utc::now();
1269
1270 for step in steps {
1271 if step.status.state.is_terminal() {
1272 continue;
1273 }
1274
1275 let (target_status, error) = match step.status.state {
1276 StepStatus::Running | StepStatus::AwaitingApproval => {
1277 let err = if step.error.is_some() {
1278 None
1279 } else {
1280 Some(error_message.to_string())
1281 };
1282 (StepStatus::Failed, err)
1283 }
1284 StepStatus::Pending => (StepStatus::Skipped, None),
1285 _ => continue,
1286 };
1287
1288 if let Err(e) = self
1289 .store
1290 .update_step(
1291 step.id,
1292 StepUpdate {
1293 status: Some(target_status),
1294 error,
1295 completed_at: Some(now),
1296 ..StepUpdate::default()
1297 },
1298 )
1299 .await
1300 {
1301 warn!(
1302 run_id = %run_id,
1303 step_id = %step.id,
1304 step_name = %step.name,
1305 error = %e,
1306 "failed to cleanup orphaned step"
1307 );
1308 } else {
1309 info!(
1310 run_id = %run_id,
1311 step_id = %step.id,
1312 step_name = %step.name,
1313 from = %step.status.state,
1314 to = %target_status,
1315 "cleaned up orphaned step"
1316 );
1317 }
1318 }
1319
1320 Ok(())
1321 }
1322
1323 async fn finalize_run(
1329 &self,
1330 run_id: Uuid,
1331 workflow_name: &str,
1332 result: Result<(), EngineError>,
1333 ctx: &WorkflowContext,
1334 run_start: Instant,
1335 run_labels: HashMap<String, String>,
1336 ) -> Result<WorkflowResult, EngineError> {
1337 let total_duration = ctx.carried_duration_ms() + run_start.elapsed().as_millis() as u64;
1340 let completed_at = Utc::now();
1341
1342 let final_status;
1343 let final_run;
1344
1345 match result {
1346 Ok(()) => {
1347 final_status = if ctx.has_allowed_failure() {
1348 RunStatus::Warning
1349 } else {
1350 RunStatus::Completed
1351 };
1352 final_run = self
1353 .store
1354 .update_run_returning(
1355 run_id,
1356 RunUpdate {
1357 status: Some(final_status),
1358 cost_usd: Some(ctx.total_cost_usd()),
1359 duration_ms: Some(total_duration),
1360 completed_at: Some(completed_at),
1361 ..RunUpdate::default()
1362 },
1363 )
1364 .await?;
1365
1366 info!(
1367 run_id = %run_id,
1368 status = %final_status,
1369 cost_usd = %ctx.total_cost_usd(),
1370 duration_ms = total_duration,
1371 "run completed"
1372 );
1373 }
1374 Err(EngineError::ApprovalRequired {
1375 run_id: approval_run_id,
1376 step_id,
1377 ref message,
1378 }) => {
1379 final_status = RunStatus::AwaitingApproval;
1380 final_run = self
1381 .store
1382 .update_run_returning(
1383 run_id,
1384 RunUpdate {
1385 status: Some(RunStatus::AwaitingApproval),
1386 cost_usd: Some(ctx.total_cost_usd()),
1387 duration_ms: Some(total_duration),
1388 ..RunUpdate::default()
1389 },
1390 )
1391 .await?;
1392
1393 info!(
1394 run_id = %approval_run_id,
1395 step_id = %step_id,
1396 message = %message,
1397 "run awaiting approval"
1398 );
1399 }
1400 Err(EngineError::DelaySleeping {
1401 run_id: delay_run_id,
1402 step_id,
1403 wake_at,
1404 }) => {
1405 final_status = RunStatus::Sleeping;
1406 final_run = self
1407 .store
1408 .update_run_returning(
1409 run_id,
1410 RunUpdate {
1411 status: Some(RunStatus::Sleeping),
1412 cost_usd: Some(ctx.total_cost_usd()),
1413 duration_ms: Some(total_duration),
1414 scheduled_at: Some(wake_at),
1415 ..RunUpdate::default()
1416 },
1417 )
1418 .await?;
1419
1420 info!(
1421 run_id = %delay_run_id,
1422 step_id = %step_id,
1423 wake_at = %wake_at,
1424 "run sleeping until delay elapses"
1425 );
1426 }
1427 Err(err) => {
1428 let guardrail_stop = matches!(
1432 err,
1433 EngineError::RunBudgetExceeded { .. } | EngineError::WorkflowGuardRejected(_)
1434 );
1435
1436 final_status = if guardrail_stop {
1437 if let Err(store_err) = self
1438 .store
1439 .update_run(
1440 run_id,
1441 RunUpdate {
1442 status: Some(RunStatus::Cancelled),
1443 error: Some(err.to_string()),
1444 cost_usd: Some(ctx.total_cost_usd()),
1445 duration_ms: Some(total_duration),
1446 completed_at: Some(completed_at),
1447 ..RunUpdate::default()
1448 },
1449 )
1450 .await
1451 {
1452 error!(run_id = %run_id, store_error = %store_err, "failed to persist run cancellation");
1453 }
1454 if let Err(cleanup_err) = self
1455 .fail_orphaned_steps(run_id, "run stopped: guardrail limit reached")
1456 .await
1457 {
1458 error!(run_id = %run_id, store_error = %cleanup_err, "failed to cleanup orphaned steps");
1459 }
1460 RunStatus::Cancelled
1461 } else {
1462 self.fail_or_schedule_retry(
1463 run_id,
1464 &err.to_string(),
1465 is_run_retryable(&err),
1466 Some(ctx.total_cost_usd()),
1467 Some(total_duration),
1468 )
1469 .await
1470 .unwrap_or_else(|store_err| {
1471 error!(run_id = %run_id, store_error = %store_err, "failed to persist run failure");
1472 RunStatus::Failed
1473 })
1474 };
1475
1476 if matches!(err, EngineError::RunBudgetExceeded { .. }) {
1477 self.on_run_budget_exceeded(workflow_name, run_id, &err);
1478 }
1479
1480 error!(run_id = %run_id, status = %final_status, error = %err, "run stopped");
1481
1482 self.publish_run_status_changed(
1483 workflow_name,
1484 run_id,
1485 final_status,
1486 Some(err.to_string()),
1487 ctx,
1488 total_duration,
1489 run_labels,
1490 );
1491
1492 #[cfg(feature = "prometheus")]
1493 self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
1494
1495 return Err(err);
1496 }
1497 }
1498
1499 self.publish_run_status_changed(
1500 workflow_name,
1501 run_id,
1502 final_status,
1503 None,
1504 ctx,
1505 total_duration,
1506 run_labels,
1507 );
1508
1509 #[cfg(feature = "prometheus")]
1510 self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
1511
1512 Ok(WorkflowResult {
1513 run: final_run,
1514 steps: ctx.step_results().to_vec(),
1515 })
1516 }
1517
1518 #[cfg(feature = "prometheus")]
1520 fn emit_run_metrics(
1521 &self,
1522 workflow_name: &str,
1523 status: RunStatus,
1524 duration_ms: u64,
1525 ctx: &WorkflowContext,
1526 ) {
1527 let status_str = status.to_string();
1528 let wf = workflow_name.to_string();
1529
1530 counter!(RUNS_TOTAL, "workflow" => wf.clone(), "status" => status_str.clone()).increment(1);
1531 histogram!(RUN_DURATION_SECONDS, "workflow" => wf.clone(), "status" => status_str)
1532 .record(duration_ms as f64 / 1000.0);
1533 histogram!(RUN_COST_USD, "workflow" => wf.clone()).record(
1534 ctx.total_cost_usd()
1535 .to_string()
1536 .parse::<f64>()
1537 .unwrap_or(0.0),
1538 );
1539 gauge!(RUNS_ACTIVE, "workflow" => wf).decrement(1.0);
1540 }
1541
1542 fn on_run_budget_exceeded(&self, workflow_name: &str, run_id: Uuid, err: &EngineError) {
1548 let EngineError::RunBudgetExceeded {
1549 limit_usd,
1550 spent_usd,
1551 step_budget_usd,
1552 ..
1553 } = err
1554 else {
1555 return;
1556 };
1557
1558 #[cfg(feature = "prometheus")]
1559 counter!(
1560 RUN_BUDGET_EXCEEDED_TOTAL,
1561 "workflow" => workflow_name.to_string(),
1562 "scope" => "run",
1563 )
1564 .increment(1);
1565
1566 self.event_publisher
1567 .publish(Event::RunBudgetExceeded(RunBudgetExceededEvent {
1568 run_id,
1569 workflow_name: workflow_name.to_string(),
1570 limit_usd: *limit_usd,
1571 spent_usd: *spent_usd,
1572 step_budget_usd: *step_budget_usd,
1573 at: Utc::now(),
1574 }));
1575 }
1576
1577 #[allow(clippy::too_many_arguments)]
1582 fn publish_run_status_changed(
1583 &self,
1584 workflow_name: &str,
1585 run_id: Uuid,
1586 to: RunStatus,
1587 error: Option<String>,
1588 ctx: &WorkflowContext,
1589 duration_ms: u64,
1590 labels: HashMap<String, String>,
1591 ) {
1592 let now = Utc::now();
1593 let cost_usd = ctx.total_cost_usd();
1594 let wf = workflow_name.to_string();
1595
1596 self.event_publisher
1597 .publish(Event::RunStatusChanged(RunStatusChangedEvent {
1598 run_id,
1599 workflow_name: wf.clone(),
1600 from: RunStatus::Running,
1601 to,
1602 error: error.clone(),
1603 cost_usd,
1604 duration_ms,
1605 labels: labels.clone(),
1606 at: now,
1607 }));
1608
1609 if to == RunStatus::Failed {
1610 self.event_publisher
1611 .publish(Event::RunFailed(RunFailedEvent {
1612 run_id,
1613 workflow_name: wf,
1614 error,
1615 cost_usd,
1616 duration_ms,
1617 labels,
1618 at: now,
1619 }));
1620 }
1621 }
1622}
1623
1624impl fmt::Debug for Engine {
1625 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1626 f.debug_struct("Engine")
1627 .field("handlers", &self.handlers.keys().collect::<Vec<_>>())
1628 .finish_non_exhaustive()
1629 }
1630}
1631
1632#[cfg(test)]
1633mod tests {
1634 use super::*;
1635 use crate::config::ShellConfig;
1636 use crate::handler::{HandlerFuture, WorkflowHandler};
1637 use ironflow_core::providers::claude::ClaudeCodeProvider;
1638 use ironflow_core::providers::record_replay::RecordReplayProvider;
1639 use ironflow_store::memory::InMemoryStore;
1640 use ironflow_store::models::StepStatus;
1641 use serde_json::json;
1642
1643 struct EchoWorkflow;
1645
1646 impl WorkflowHandler for EchoWorkflow {
1647 fn name(&self) -> &str {
1648 "echo-workflow"
1649 }
1650
1651 fn describe(&self) -> WorkflowInfo {
1652 WorkflowInfo {
1653 description: "A simple workflow that echoes hello".to_string(),
1654 source_code: None,
1655 sub_workflows: Vec::new(),
1656 category: None,
1657 version: self.version().map(str::to_string),
1658 compatible_versions: Vec::new(),
1659 input_schema: None,
1660 default_labels: HashMap::new(),
1661 schedule: self.schedule().cloned(),
1662 default_max_cost_usd: self.default_max_cost_usd(),
1663 }
1664 }
1665
1666 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
1667 Box::pin(async move {
1668 ctx.shell("greet", ShellConfig::new("echo hello")).await?;
1669 Ok(())
1670 })
1671 }
1672 }
1673
1674 struct FailingWorkflow;
1676
1677 impl WorkflowHandler for FailingWorkflow {
1678 fn name(&self) -> &str {
1679 "failing-workflow"
1680 }
1681
1682 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
1683 Box::pin(async move {
1684 ctx.shell("fail", ShellConfig::new("exit 1")).await?;
1685 Ok(())
1686 })
1687 }
1688 }
1689
1690 fn create_test_engine() -> Engine {
1691 let store = Arc::new(InMemoryStore::new());
1692 let inner = ClaudeCodeProvider::new();
1693 let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
1694 inner,
1695 "/tmp/ironflow-fixtures",
1696 ));
1697 Engine::new(store, provider)
1698 }
1699
1700 #[test]
1701 fn engine_new_creates_instance() {
1702 let engine = create_test_engine();
1703 assert_eq!(engine.handler_names().len(), 0);
1704 }
1705
1706 #[test]
1707 fn engine_register_handler() {
1708 let mut engine = create_test_engine();
1709 let result = engine.register(EchoWorkflow);
1710 assert!(result.is_ok());
1711 assert_eq!(engine.handler_names().len(), 1);
1712 assert!(engine.handler_names().contains(&"echo-workflow"));
1713 }
1714
1715 #[test]
1716 fn engine_register_duplicate_returns_error() {
1717 let mut engine = create_test_engine();
1718 engine.register(EchoWorkflow).unwrap();
1719 let result = engine.register(EchoWorkflow);
1720 assert!(result.is_err());
1721 }
1722
1723 #[test]
1724 fn engine_get_handler_found() {
1725 let mut engine = create_test_engine();
1726 engine.register(EchoWorkflow).unwrap();
1727 let handler = engine.get_handler("echo-workflow");
1728 assert!(handler.is_some());
1729 }
1730
1731 #[test]
1732 fn engine_get_handler_not_found() {
1733 let engine = create_test_engine();
1734 let handler = engine.get_handler("nonexistent");
1735 assert!(handler.is_none());
1736 }
1737
1738 #[test]
1739 fn engine_handler_names_lists_all() {
1740 let mut engine = create_test_engine();
1741 engine.register(EchoWorkflow).unwrap();
1742 engine.register(FailingWorkflow).unwrap();
1743 let names = engine.handler_names();
1744 assert_eq!(names.len(), 2);
1745 assert!(names.contains(&"echo-workflow"));
1746 assert!(names.contains(&"failing-workflow"));
1747 }
1748
1749 #[test]
1750 fn engine_handler_info_returns_description() {
1751 let mut engine = create_test_engine();
1752 engine.register(EchoWorkflow).unwrap();
1753 let info = engine.handler_info("echo-workflow");
1754 assert!(info.is_some());
1755 let info = info.unwrap();
1756 assert_eq!(info.description, "A simple workflow that echoes hello");
1757 }
1758
1759 struct CategorizedWorkflow;
1760
1761 impl WorkflowHandler for CategorizedWorkflow {
1762 fn name(&self) -> &str {
1763 "categorized"
1764 }
1765 fn category(&self) -> Option<&str> {
1766 Some("data/etl")
1767 }
1768 fn execute<'a>(
1769 &'a self,
1770 _ctx: &'a mut WorkflowContext,
1771 ) -> crate::handler::HandlerFuture<'a> {
1772 Box::pin(async move { Ok(()) })
1773 }
1774 }
1775
1776 #[test]
1777 fn engine_default_describe_propagates_category() {
1778 let mut engine = create_test_engine();
1779 engine.register(CategorizedWorkflow).unwrap();
1780 let info = engine.handler_info("categorized").unwrap();
1781 assert_eq!(info.category.as_deref(), Some("data/etl"));
1782 }
1783
1784 #[test]
1785 fn engine_default_describe_without_category() {
1786 let mut engine = create_test_engine();
1787 engine.register(EchoWorkflow).unwrap();
1788 let info = engine.handler_info("echo-workflow").unwrap();
1789 assert!(info.category.is_none());
1790 }
1791
1792 struct ScheduledWorkflow {
1797 schedule: CronSchedule,
1798 }
1799
1800 impl ScheduledWorkflow {
1801 fn new() -> Self {
1802 Self {
1803 schedule: CronSchedule::new("0 0 * * * *").unwrap(),
1804 }
1805 }
1806 }
1807
1808 impl WorkflowHandler for ScheduledWorkflow {
1809 fn name(&self) -> &str {
1810 "scheduled"
1811 }
1812 fn schedule(&self) -> Option<&CronSchedule> {
1813 Some(&self.schedule)
1814 }
1815 fn execute<'a>(
1816 &'a self,
1817 _ctx: &'a mut WorkflowContext,
1818 ) -> crate::handler::HandlerFuture<'a> {
1819 Box::pin(async move { Ok(()) })
1820 }
1821 }
1822
1823 #[test]
1824 fn engine_default_describe_propagates_schedule() {
1825 let mut engine = create_test_engine();
1826 engine.register(ScheduledWorkflow::new()).unwrap();
1827 let info = engine.handler_info("scheduled").unwrap();
1828 assert_eq!(
1829 info.schedule.as_ref().map(|s| s.as_str()),
1830 Some("0 0 * * * *")
1831 );
1832 }
1833
1834 #[test]
1835 fn engine_default_describe_without_schedule() {
1836 let mut engine = create_test_engine();
1837 engine.register(EchoWorkflow).unwrap();
1838 let info = engine.handler_info("echo-workflow").unwrap();
1839 assert!(info.schedule.is_none());
1840 }
1841
1842 #[test]
1843 fn scheduled_handlers_returns_only_scheduled() {
1844 let mut engine = create_test_engine();
1845 engine.register(EchoWorkflow).unwrap();
1846 engine.register(ScheduledWorkflow::new()).unwrap();
1847 engine.register(FailingWorkflow).unwrap();
1848
1849 let scheduled = engine.scheduled_handlers();
1850 assert_eq!(scheduled.len(), 1);
1851 assert_eq!(scheduled[0].0, "scheduled");
1852 assert_eq!(scheduled[0].1.as_str(), "0 0 * * * *");
1853 }
1854
1855 #[test]
1856 fn scheduled_handlers_empty_when_none_scheduled() {
1857 let mut engine = create_test_engine();
1858 engine.register(EchoWorkflow).unwrap();
1859 engine.register(FailingWorkflow).unwrap();
1860
1861 let scheduled = engine.scheduled_handlers();
1862 assert!(scheduled.is_empty());
1863 }
1864
1865 struct BadCategoryWorkflow(&'static str);
1866
1867 impl WorkflowHandler for BadCategoryWorkflow {
1868 fn name(&self) -> &str {
1869 "bad-category"
1870 }
1871 fn category(&self) -> Option<&str> {
1872 Some(self.0)
1873 }
1874 fn execute<'a>(
1875 &'a self,
1876 _ctx: &'a mut WorkflowContext,
1877 ) -> crate::handler::HandlerFuture<'a> {
1878 Box::pin(async move { Ok(()) })
1879 }
1880 }
1881
1882 #[test]
1883 fn engine_register_rejects_empty_category() {
1884 let mut engine = create_test_engine();
1885 let err = engine.register(BadCategoryWorkflow("")).unwrap_err();
1886 match err {
1887 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty category")),
1888 other => panic!("expected InvalidWorkflow, got {other:?}"),
1889 }
1890 }
1891
1892 #[test]
1893 fn engine_register_rejects_leading_slash_category() {
1894 let mut engine = create_test_engine();
1895 let err = engine
1896 .register(BadCategoryWorkflow("/data/etl"))
1897 .unwrap_err();
1898 match err {
1899 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("leading '/'")),
1900 other => panic!("expected InvalidWorkflow, got {other:?}"),
1901 }
1902 }
1903
1904 #[test]
1905 fn engine_register_rejects_trailing_slash_category() {
1906 let mut engine = create_test_engine();
1907 let err = engine
1908 .register(BadCategoryWorkflow("data/etl/"))
1909 .unwrap_err();
1910 match err {
1911 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("trailing '/'")),
1912 other => panic!("expected InvalidWorkflow, got {other:?}"),
1913 }
1914 }
1915
1916 #[test]
1917 fn engine_register_rejects_double_slash_category() {
1918 let mut engine = create_test_engine();
1919 let err = engine
1920 .register(BadCategoryWorkflow("data//etl"))
1921 .unwrap_err();
1922 match err {
1923 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty segment")),
1924 other => panic!("expected InvalidWorkflow, got {other:?}"),
1925 }
1926 }
1927
1928 #[test]
1929 fn engine_register_rejects_whitespace_only_segment_category() {
1930 let mut engine = create_test_engine();
1931 let err = engine
1932 .register(BadCategoryWorkflow("data/ /etl"))
1933 .unwrap_err();
1934 match err {
1935 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("whitespace-only segment")),
1936 other => panic!("expected InvalidWorkflow, got {other:?}"),
1937 }
1938 }
1939
1940 #[test]
1941 fn engine_register_accepts_valid_nested_category() {
1942 let mut engine = create_test_engine();
1943 assert!(engine.register(CategorizedWorkflow).is_ok());
1944 }
1945
1946 #[tokio::test]
1947 async fn engine_unknown_workflow_returns_error() {
1948 let engine = create_test_engine();
1949 let result = engine
1950 .run_handler("unknown", TriggerKind::Manual, json!({}))
1951 .await;
1952 assert!(result.is_err());
1953 match result {
1954 Err(EngineError::InvalidWorkflow(msg)) => {
1955 assert!(msg.contains("no handler registered"));
1956 }
1957 _ => panic!("expected InvalidWorkflow error"),
1958 }
1959 }
1960
1961 #[tokio::test]
1962 async fn engine_enqueue_handler_creates_pending_run() {
1963 let mut engine = create_test_engine();
1964 engine.register(EchoWorkflow).unwrap();
1965
1966 let run = engine
1967 .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
1968 .await
1969 .unwrap();
1970 assert_eq!(run.status.state, RunStatus::Pending);
1971 assert_eq!(run.workflow_name, "echo-workflow");
1972 }
1973
1974 #[tokio::test]
1975 async fn enqueue_handler_leaves_the_run_unattributed() {
1976 let mut engine = create_test_engine();
1977 engine.register(EchoWorkflow).unwrap();
1978
1979 let run = engine
1980 .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
1981 .await
1982 .unwrap();
1983
1984 assert!(run.created_by.is_none());
1985 }
1986
1987 #[tokio::test]
1988 async fn enqueue_handler_with_options_records_the_author() {
1989 let mut engine = create_test_engine();
1990 engine.register(EchoWorkflow).unwrap();
1991 let actor = RunActor::User {
1992 user_id: Uuid::now_v7(),
1993 };
1994
1995 let run = engine
1996 .enqueue_handler_with_options(
1997 "echo-workflow",
1998 TriggerKind::Api,
1999 json!({}),
2000 EnqueueOptions {
2001 created_by: Some(actor.clone()),
2002 ..Default::default()
2003 },
2004 )
2005 .await
2006 .unwrap()
2007 .into_run();
2008
2009 assert_eq!(run.created_by, Some(actor));
2010 }
2011
2012 #[tokio::test]
2013 async fn enqueue_handler_with_options_accepts_no_author() {
2014 let mut engine = create_test_engine();
2015 engine.register(EchoWorkflow).unwrap();
2016
2017 let run = engine
2018 .enqueue_handler_with_options(
2019 "echo-workflow",
2020 TriggerKind::Cron {
2021 schedule: "0 * * * * *".to_string(),
2022 },
2023 json!({}),
2024 EnqueueOptions::default(),
2025 )
2026 .await
2027 .unwrap()
2028 .into_run();
2029
2030 assert!(run.created_by.is_none());
2031 }
2032
2033 #[tokio::test]
2034 async fn run_handler_leaves_the_run_unattributed() {
2035 let mut engine = create_test_engine();
2036 engine.register(EchoWorkflow).unwrap();
2037
2038 let run = engine
2039 .run_handler("echo-workflow", TriggerKind::Manual, json!({}))
2040 .await
2041 .unwrap()
2042 .run;
2043
2044 assert!(run.created_by.is_none());
2045 }
2046
2047 #[tokio::test]
2048 async fn engine_register_boxed() {
2049 let mut engine = create_test_engine();
2050 let handler: Box<dyn WorkflowHandler> = Box::new(EchoWorkflow);
2051 let result = engine.register_boxed(handler);
2052 assert!(result.is_ok());
2053 assert_eq!(engine.handler_names().len(), 1);
2054 }
2055
2056 #[tokio::test]
2057 async fn engine_store_and_provider_accessors() {
2058 let store = Arc::new(InMemoryStore::new());
2059 let inner = ClaudeCodeProvider::new();
2060 let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
2061 inner,
2062 "/tmp/ironflow-fixtures",
2063 ));
2064 let engine = Engine::new(store.clone(), provider.clone());
2065
2066 let _ = engine.store();
2068 let _ = engine.provider();
2069 }
2070
2071 use crate::operation::{Operation, OperationContext};
2076 use async_trait::async_trait;
2077 use ironflow_core::error::OperationError;
2078 use ironflow_store::models::StepKind;
2079
2080 struct FakeGitlabOp {
2081 project_id: u64,
2082 title: String,
2083 }
2084
2085 #[async_trait]
2086 impl Operation for FakeGitlabOp {
2087 fn kind(&self) -> &str {
2088 "gitlab"
2089 }
2090
2091 async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
2092 Ok(json!({
2093 "issue_id": 42,
2094 "project_id": self.project_id,
2095 "title": self.title,
2096 }))
2097 }
2098
2099 fn input(&self) -> Option<Value> {
2100 Some(json!({
2101 "project_id": self.project_id,
2102 "title": self.title,
2103 }))
2104 }
2105 }
2106
2107 struct FailingOp;
2108
2109 #[async_trait]
2110 impl Operation for FailingOp {
2111 fn kind(&self) -> &str {
2112 "broken-service"
2113 }
2114
2115 async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
2116 Err(OperationError::Http {
2117 status: None,
2118 message: "service unavailable".to_string(),
2119 })
2120 }
2121 }
2122
2123 struct OperationWorkflow;
2124
2125 impl WorkflowHandler for OperationWorkflow {
2126 fn name(&self) -> &str {
2127 "operation-workflow"
2128 }
2129
2130 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2131 Box::pin(async move {
2132 let op = FakeGitlabOp {
2133 project_id: 123,
2134 title: "Bug report".to_string(),
2135 };
2136 ctx.operation("create-issue", &op).await?;
2137 Ok(())
2138 })
2139 }
2140 }
2141
2142 struct FailingOperationWorkflow;
2143
2144 impl WorkflowHandler for FailingOperationWorkflow {
2145 fn name(&self) -> &str {
2146 "failing-operation-workflow"
2147 }
2148
2149 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2150 Box::pin(async move {
2151 ctx.operation("broken-call", &FailingOp).await?;
2152 Ok(())
2153 })
2154 }
2155 }
2156
2157 struct MixedWorkflow;
2158
2159 impl WorkflowHandler for MixedWorkflow {
2160 fn name(&self) -> &str {
2161 "mixed-workflow"
2162 }
2163
2164 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2165 Box::pin(async move {
2166 ctx.shell("build", ShellConfig::new("echo built")).await?;
2167 let op = FakeGitlabOp {
2168 project_id: 456,
2169 title: "Deploy done".to_string(),
2170 };
2171 let result = ctx.operation("notify-gitlab", &op).await?;
2172 assert_eq!(result.output["issue_id"], 42);
2173 Ok(())
2174 })
2175 }
2176 }
2177
2178 #[tokio::test]
2179 async fn operation_step_happy_path() {
2180 let mut engine = create_test_engine();
2181 engine.register(OperationWorkflow).unwrap();
2182
2183 let run = engine
2184 .run_handler("operation-workflow", TriggerKind::Manual, json!({}))
2185 .await
2186 .unwrap()
2187 .run;
2188
2189 assert_eq!(run.status.state, RunStatus::Completed);
2190
2191 let steps = engine.store().list_steps(run.id).await.unwrap();
2192
2193 assert_eq!(steps.len(), 1);
2194 assert_eq!(steps[0].name, "create-issue");
2195 assert_eq!(steps[0].kind, StepKind::Custom("gitlab".to_string()));
2196 assert_eq!(
2197 steps[0].status.state,
2198 ironflow_store::models::StepStatus::Completed
2199 );
2200
2201 let output = steps[0].output.as_ref().unwrap();
2202 assert_eq!(output["issue_id"], 42);
2203 assert_eq!(output["project_id"], 123);
2204
2205 let input = steps[0].input.as_ref().unwrap();
2206 assert_eq!(input["project_id"], 123);
2207 assert_eq!(input["title"], "Bug report");
2208 }
2209
2210 #[tokio::test]
2211 async fn operation_step_failure_marks_run_failed() {
2212 let mut engine = create_test_engine();
2213 engine.register(FailingOperationWorkflow).unwrap();
2214
2215 let result = engine
2216 .run_handler("failing-operation-workflow", TriggerKind::Manual, json!({}))
2217 .await;
2218
2219 assert!(result.is_err());
2220 }
2221
2222 #[tokio::test]
2223 async fn operation_mixed_with_shell_steps() {
2224 let mut engine = create_test_engine();
2225 engine.register(MixedWorkflow).unwrap();
2226
2227 let run = engine
2228 .run_handler("mixed-workflow", TriggerKind::Manual, json!({}))
2229 .await
2230 .unwrap()
2231 .run;
2232
2233 assert_eq!(run.status.state, RunStatus::Completed);
2234
2235 let steps = engine.store().list_steps(run.id).await.unwrap();
2236
2237 assert_eq!(steps.len(), 2);
2238 assert_eq!(steps[0].kind, StepKind::Shell);
2239 assert_eq!(steps[1].kind, StepKind::Custom("gitlab".to_string()));
2240 assert_eq!(steps[0].position, 0);
2241 assert_eq!(steps[1].position, 1);
2242 }
2243
2244 use crate::config::ApprovalConfig;
2249
2250 struct SingleApprovalWorkflow;
2251
2252 impl WorkflowHandler for SingleApprovalWorkflow {
2253 fn name(&self) -> &str {
2254 "single-approval"
2255 }
2256
2257 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2258 Box::pin(async move {
2259 ctx.shell("build", ShellConfig::new("echo built")).await?;
2260 ctx.approval("gate", ApprovalConfig::new("OK?")).await?;
2261 ctx.shell("deploy", ShellConfig::new("echo deployed"))
2262 .await?;
2263 Ok(())
2264 })
2265 }
2266 }
2267
2268 struct DoubleApprovalWorkflow;
2269
2270 impl WorkflowHandler for DoubleApprovalWorkflow {
2271 fn name(&self) -> &str {
2272 "double-approval"
2273 }
2274
2275 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2276 Box::pin(async move {
2277 ctx.shell("build", ShellConfig::new("echo built")).await?;
2278 ctx.approval("staging-gate", ApprovalConfig::new("Deploy staging?"))
2279 .await?;
2280 ctx.shell("deploy-staging", ShellConfig::new("echo staging"))
2281 .await?;
2282 ctx.approval("prod-gate", ApprovalConfig::new("Deploy prod?"))
2283 .await?;
2284 ctx.shell("deploy-prod", ShellConfig::new("echo prod"))
2285 .await?;
2286 Ok(())
2287 })
2288 }
2289 }
2290
2291 #[tokio::test]
2292 async fn approval_pauses_run() {
2293 let mut engine = create_test_engine();
2294 engine.register(SingleApprovalWorkflow).unwrap();
2295
2296 let run = engine
2297 .run_handler("single-approval", TriggerKind::Manual, json!({}))
2298 .await
2299 .unwrap()
2300 .run;
2301
2302 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
2303
2304 let steps = engine.store().list_steps(run.id).await.unwrap();
2305 assert_eq!(steps.len(), 2); assert_eq!(steps[0].kind, StepKind::Shell);
2307 assert_eq!(steps[0].status.state, StepStatus::Completed);
2308 assert_eq!(steps[1].kind, StepKind::Approval);
2309 assert_eq!(steps[1].status.state, StepStatus::AwaitingApproval);
2310 }
2311
2312 #[tokio::test]
2313 async fn approval_resume_completes_run() {
2314 let mut engine = create_test_engine();
2315 engine.register(SingleApprovalWorkflow).unwrap();
2316
2317 let run = engine
2319 .run_handler("single-approval", TriggerKind::Manual, json!({}))
2320 .await
2321 .unwrap()
2322 .run;
2323 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
2324
2325 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;
2334 assert_eq!(resumed.status.state, RunStatus::Completed);
2335
2336 let steps = engine.store().list_steps(run.id).await.unwrap();
2337 assert_eq!(steps.len(), 3); assert_eq!(steps[0].name, "build");
2339 assert_eq!(steps[0].status.state, StepStatus::Completed);
2340 assert_eq!(steps[1].name, "gate");
2341 assert_eq!(steps[1].kind, StepKind::Approval);
2342 assert_eq!(steps[1].status.state, StepStatus::Completed);
2343 assert_eq!(steps[2].name, "deploy");
2344 assert_eq!(steps[2].status.state, StepStatus::Completed);
2345 }
2346
2347 #[tokio::test]
2348 async fn double_approval_two_resumes() {
2349 let mut engine = create_test_engine();
2350 engine.register(DoubleApprovalWorkflow).unwrap();
2351
2352 let run = engine
2354 .run_handler("double-approval", TriggerKind::Manual, json!({}))
2355 .await
2356 .unwrap()
2357 .run;
2358 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
2359
2360 let steps = engine.store().list_steps(run.id).await.unwrap();
2361 assert_eq!(steps.len(), 2); engine
2365 .store()
2366 .update_run_status(run.id, RunStatus::Running)
2367 .await
2368 .unwrap();
2369
2370 let resumed = engine.resume_run(run.id).await.unwrap().run;
2371 assert_eq!(resumed.status.state, RunStatus::AwaitingApproval);
2372
2373 let steps = engine.store().list_steps(run.id).await.unwrap();
2374 assert_eq!(steps.len(), 4); engine
2378 .store()
2379 .update_run_status(run.id, RunStatus::Running)
2380 .await
2381 .unwrap();
2382
2383 let final_run = engine.resume_run(run.id).await.unwrap().run;
2384 assert_eq!(final_run.status.state, RunStatus::Completed);
2385
2386 let steps = engine.store().list_steps(run.id).await.unwrap();
2387 assert_eq!(steps.len(), 5);
2388 assert_eq!(steps[0].name, "build");
2389 assert_eq!(steps[1].name, "staging-gate");
2390 assert_eq!(steps[2].name, "deploy-staging");
2391 assert_eq!(steps[3].name, "prod-gate");
2392 assert_eq!(steps[4].name, "deploy-prod");
2393
2394 for step in &steps {
2395 assert_eq!(step.status.state, StepStatus::Completed);
2396 }
2397 }
2398
2399 use ironflow_store::models::{NewStep, StepUpdate, step_trace_id};
2404
2405 async fn create_step_with_status(
2406 store: &Arc<dyn Store>,
2407 run_id: Uuid,
2408 name: &str,
2409 position: u32,
2410 status: StepStatus,
2411 ) -> ironflow_store::models::Step {
2412 let step = store
2413 .create_step(NewStep {
2414 run_id,
2415 trace_id: step_trace_id(run_id, name, position),
2416 name: name.to_string(),
2417 kind: StepKind::Shell,
2418 position,
2419 input: None,
2420 is_error_handler: false,
2421 })
2422 .await
2423 .unwrap();
2424
2425 match status {
2426 StepStatus::Pending => {}
2427 StepStatus::Running => {
2428 store
2429 .update_step(
2430 step.id,
2431 StepUpdate {
2432 status: Some(StepStatus::Running),
2433 ..StepUpdate::default()
2434 },
2435 )
2436 .await
2437 .unwrap();
2438 }
2439 StepStatus::Completed => {
2440 store
2441 .update_step(
2442 step.id,
2443 StepUpdate {
2444 status: Some(StepStatus::Running),
2445 ..StepUpdate::default()
2446 },
2447 )
2448 .await
2449 .unwrap();
2450 store
2451 .update_step(
2452 step.id,
2453 StepUpdate {
2454 status: Some(StepStatus::Completed),
2455 ..StepUpdate::default()
2456 },
2457 )
2458 .await
2459 .unwrap();
2460 }
2461 StepStatus::AwaitingApproval => {
2462 store
2463 .update_step(
2464 step.id,
2465 StepUpdate {
2466 status: Some(StepStatus::Running),
2467 ..StepUpdate::default()
2468 },
2469 )
2470 .await
2471 .unwrap();
2472 store
2473 .update_step(
2474 step.id,
2475 StepUpdate {
2476 status: Some(StepStatus::AwaitingApproval),
2477 ..StepUpdate::default()
2478 },
2479 )
2480 .await
2481 .unwrap();
2482 }
2483 _ => panic!("unsupported status for test helper: {status}"),
2484 }
2485
2486 store.get_step(step.id).await.unwrap().unwrap()
2487 }
2488
2489 #[tokio::test]
2490 async fn fail_orphaned_steps_marks_running_as_failed() {
2491 let engine = create_test_engine();
2492 let run = engine
2493 .store()
2494 .create_run(NewRun {
2495 created_by: None,
2496 workflow_name: "test".to_string(),
2497 trigger: TriggerKind::Manual,
2498 payload: json!({}),
2499 max_retries: 0,
2500 handler_version: None,
2501 labels: HashMap::new(),
2502 scheduled_at: None,
2503 idempotency_key: None,
2504 max_cost_usd: None,
2505 })
2506 .await
2507 .unwrap()
2508 .into_run();
2509
2510 let step = create_step_with_status(
2511 engine.store(),
2512 run.id,
2513 "running-step",
2514 0,
2515 StepStatus::Running,
2516 )
2517 .await;
2518
2519 engine
2520 .fail_orphaned_steps(run.id, "parent run timed out")
2521 .await
2522 .unwrap();
2523
2524 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
2525 assert_eq!(updated.status.state, StepStatus::Failed);
2526 assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
2527 assert!(updated.completed_at.is_some());
2528 }
2529
2530 #[tokio::test]
2531 async fn fail_orphaned_steps_marks_pending_as_skipped() {
2532 let engine = create_test_engine();
2533 let run = engine
2534 .store()
2535 .create_run(NewRun {
2536 created_by: None,
2537 workflow_name: "test".to_string(),
2538 trigger: TriggerKind::Manual,
2539 payload: json!({}),
2540 max_retries: 0,
2541 handler_version: None,
2542 labels: HashMap::new(),
2543 scheduled_at: None,
2544 idempotency_key: None,
2545 max_cost_usd: None,
2546 })
2547 .await
2548 .unwrap()
2549 .into_run();
2550
2551 let step = create_step_with_status(
2552 engine.store(),
2553 run.id,
2554 "pending-step",
2555 0,
2556 StepStatus::Pending,
2557 )
2558 .await;
2559
2560 engine
2561 .fail_orphaned_steps(run.id, "parent run timed out")
2562 .await
2563 .unwrap();
2564
2565 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
2566 assert_eq!(updated.status.state, StepStatus::Skipped);
2567 assert!(updated.error.is_none());
2568 assert!(updated.completed_at.is_some());
2569 }
2570
2571 #[tokio::test]
2572 async fn fail_orphaned_steps_marks_awaiting_approval_as_failed() {
2573 let engine = create_test_engine();
2574 let run = engine
2575 .store()
2576 .create_run(NewRun {
2577 created_by: None,
2578 workflow_name: "test".to_string(),
2579 trigger: TriggerKind::Manual,
2580 payload: json!({}),
2581 max_retries: 0,
2582 handler_version: None,
2583 labels: HashMap::new(),
2584 scheduled_at: None,
2585 idempotency_key: None,
2586 max_cost_usd: None,
2587 })
2588 .await
2589 .unwrap()
2590 .into_run();
2591
2592 let step = create_step_with_status(
2593 engine.store(),
2594 run.id,
2595 "approval-step",
2596 0,
2597 StepStatus::AwaitingApproval,
2598 )
2599 .await;
2600
2601 engine
2602 .fail_orphaned_steps(run.id, "parent run timed out")
2603 .await
2604 .unwrap();
2605
2606 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
2607 assert_eq!(updated.status.state, StepStatus::Failed);
2608 assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
2609 assert!(updated.completed_at.is_some());
2610 }
2611
2612 #[tokio::test]
2613 async fn fail_orphaned_steps_skips_terminal_steps() {
2614 let engine = create_test_engine();
2615 let run = engine
2616 .store()
2617 .create_run(NewRun {
2618 created_by: None,
2619 workflow_name: "test".to_string(),
2620 trigger: TriggerKind::Manual,
2621 payload: json!({}),
2622 max_retries: 0,
2623 handler_version: None,
2624 labels: HashMap::new(),
2625 scheduled_at: None,
2626 idempotency_key: None,
2627 max_cost_usd: None,
2628 })
2629 .await
2630 .unwrap()
2631 .into_run();
2632
2633 let completed_step =
2634 create_step_with_status(engine.store(), run.id, "done", 0, StepStatus::Completed).await;
2635 let running_step =
2636 create_step_with_status(engine.store(), run.id, "in-flight", 1, StepStatus::Running)
2637 .await;
2638
2639 engine
2640 .fail_orphaned_steps(run.id, "parent run timed out")
2641 .await
2642 .unwrap();
2643
2644 let completed = engine
2645 .store()
2646 .get_step(completed_step.id)
2647 .await
2648 .unwrap()
2649 .unwrap();
2650 assert_eq!(completed.status.state, StepStatus::Completed);
2651
2652 let failed = engine
2653 .store()
2654 .get_step(running_step.id)
2655 .await
2656 .unwrap()
2657 .unwrap();
2658 assert_eq!(failed.status.state, StepStatus::Failed);
2659 }
2660
2661 #[tokio::test]
2662 async fn fail_orphaned_steps_mixed_states() {
2663 let engine = create_test_engine();
2664 let run = engine
2665 .store()
2666 .create_run(NewRun {
2667 created_by: None,
2668 workflow_name: "test".to_string(),
2669 trigger: TriggerKind::Manual,
2670 payload: json!({}),
2671 max_retries: 0,
2672 handler_version: None,
2673 labels: HashMap::new(),
2674 scheduled_at: None,
2675 idempotency_key: None,
2676 max_cost_usd: None,
2677 })
2678 .await
2679 .unwrap()
2680 .into_run();
2681
2682 let s_completed =
2683 create_step_with_status(engine.store(), run.id, "step-1", 0, StepStatus::Completed)
2684 .await;
2685 let s_running =
2686 create_step_with_status(engine.store(), run.id, "step-2", 1, StepStatus::Running).await;
2687 let s_pending =
2688 create_step_with_status(engine.store(), run.id, "step-3", 2, StepStatus::Pending).await;
2689
2690 engine.fail_orphaned_steps(run.id, "timeout").await.unwrap();
2691
2692 let r_completed = engine
2693 .store()
2694 .get_step(s_completed.id)
2695 .await
2696 .unwrap()
2697 .unwrap();
2698 assert_eq!(r_completed.status.state, StepStatus::Completed);
2699
2700 let r_running = engine
2701 .store()
2702 .get_step(s_running.id)
2703 .await
2704 .unwrap()
2705 .unwrap();
2706 assert_eq!(r_running.status.state, StepStatus::Failed);
2707 assert_eq!(r_running.error.as_deref(), Some("timeout"));
2708
2709 let r_pending = engine
2710 .store()
2711 .get_step(s_pending.id)
2712 .await
2713 .unwrap()
2714 .unwrap();
2715 assert_eq!(r_pending.status.state, StepStatus::Skipped);
2716 assert!(r_pending.error.is_none());
2717 }
2718
2719 #[tokio::test]
2720 async fn fail_orphaned_steps_no_steps_is_noop() {
2721 let engine = create_test_engine();
2722 let run = engine
2723 .store()
2724 .create_run(NewRun {
2725 created_by: None,
2726 workflow_name: "test".to_string(),
2727 trigger: TriggerKind::Manual,
2728 payload: json!({}),
2729 max_retries: 0,
2730 handler_version: None,
2731 labels: HashMap::new(),
2732 scheduled_at: None,
2733 idempotency_key: None,
2734 max_cost_usd: None,
2735 })
2736 .await
2737 .unwrap()
2738 .into_run();
2739
2740 let result = engine.fail_orphaned_steps(run.id, "timeout").await;
2741 assert!(result.is_ok());
2742 }
2743
2744 #[tokio::test]
2745 async fn fail_orphaned_steps_preserves_existing_error() {
2746 let engine = create_test_engine();
2747 let run = engine
2748 .store()
2749 .create_run(NewRun {
2750 created_by: None,
2751 workflow_name: "test".to_string(),
2752 trigger: TriggerKind::Manual,
2753 payload: json!({}),
2754 max_retries: 0,
2755 handler_version: None,
2756 labels: HashMap::new(),
2757 scheduled_at: None,
2758 idempotency_key: None,
2759 max_cost_usd: None,
2760 })
2761 .await
2762 .unwrap()
2763 .into_run();
2764
2765 let step_with_error = create_step_with_status(
2766 engine.store(),
2767 run.id,
2768 "already-errored",
2769 0,
2770 StepStatus::Running,
2771 )
2772 .await;
2773
2774 engine
2775 .store()
2776 .update_step(
2777 step_with_error.id,
2778 StepUpdate {
2779 error: Some("real error from provider".to_string()),
2780 ..StepUpdate::default()
2781 },
2782 )
2783 .await
2784 .unwrap();
2785
2786 let step_no_error = create_step_with_status(
2787 engine.store(),
2788 run.id,
2789 "no-error-yet",
2790 1,
2791 StepStatus::Running,
2792 )
2793 .await;
2794
2795 engine
2796 .fail_orphaned_steps(run.id, "parent run failed")
2797 .await
2798 .unwrap();
2799
2800 let updated_with = engine
2801 .store()
2802 .get_step(step_with_error.id)
2803 .await
2804 .unwrap()
2805 .unwrap();
2806 assert_eq!(updated_with.status.state, StepStatus::Failed);
2807 assert_eq!(
2808 updated_with.error.as_deref(),
2809 Some("real error from provider"),
2810 );
2811
2812 let updated_without = engine
2813 .store()
2814 .get_step(step_no_error.id)
2815 .await
2816 .unwrap()
2817 .unwrap();
2818 assert_eq!(updated_without.status.state, StepStatus::Failed);
2819 assert_eq!(updated_without.error.as_deref(), Some("parent run failed"),);
2820 }
2821}