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 ApprovalRequestedEvent, Event, EventPublisher, EventSubscriber, RunBudgetExceededEvent,
44 RunFailedEvent, 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 let requirement = self
1402 .store
1403 .get_step(step_id)
1404 .await?
1405 .and_then(|s| s.approval_requirement);
1406 self.event_publisher
1407 .publish(Event::ApprovalRequested(ApprovalRequestedEvent {
1408 run_id: approval_run_id,
1409 step_id,
1410 message: message.clone(),
1411 requirement,
1412 at: Utc::now(),
1413 }));
1414 }
1415 Err(EngineError::DelaySleeping {
1416 run_id: delay_run_id,
1417 step_id,
1418 wake_at,
1419 }) => {
1420 final_status = RunStatus::Sleeping;
1421 final_run = self
1422 .store
1423 .update_run_returning(
1424 run_id,
1425 RunUpdate {
1426 status: Some(RunStatus::Sleeping),
1427 cost_usd: Some(ctx.total_cost_usd()),
1428 duration_ms: Some(total_duration),
1429 scheduled_at: Some(wake_at),
1430 ..RunUpdate::default()
1431 },
1432 )
1433 .await?;
1434
1435 info!(
1436 run_id = %delay_run_id,
1437 step_id = %step_id,
1438 wake_at = %wake_at,
1439 "run sleeping until delay elapses"
1440 );
1441 }
1442 Err(err) => {
1443 let guardrail_stop = matches!(
1447 err,
1448 EngineError::RunBudgetExceeded { .. } | EngineError::WorkflowGuardRejected(_)
1449 );
1450
1451 final_status = if guardrail_stop {
1452 if let Err(store_err) = self
1453 .store
1454 .update_run(
1455 run_id,
1456 RunUpdate {
1457 status: Some(RunStatus::Cancelled),
1458 error: Some(err.to_string()),
1459 cost_usd: Some(ctx.total_cost_usd()),
1460 duration_ms: Some(total_duration),
1461 completed_at: Some(completed_at),
1462 ..RunUpdate::default()
1463 },
1464 )
1465 .await
1466 {
1467 error!(run_id = %run_id, store_error = %store_err, "failed to persist run cancellation");
1468 }
1469 if let Err(cleanup_err) = self
1470 .fail_orphaned_steps(run_id, "run stopped: guardrail limit reached")
1471 .await
1472 {
1473 error!(run_id = %run_id, store_error = %cleanup_err, "failed to cleanup orphaned steps");
1474 }
1475 RunStatus::Cancelled
1476 } else {
1477 self.fail_or_schedule_retry(
1478 run_id,
1479 &err.to_string(),
1480 is_run_retryable(&err),
1481 Some(ctx.total_cost_usd()),
1482 Some(total_duration),
1483 )
1484 .await
1485 .unwrap_or_else(|store_err| {
1486 error!(run_id = %run_id, store_error = %store_err, "failed to persist run failure");
1487 RunStatus::Failed
1488 })
1489 };
1490
1491 if matches!(err, EngineError::RunBudgetExceeded { .. }) {
1492 self.on_run_budget_exceeded(workflow_name, run_id, &err);
1493 }
1494
1495 error!(run_id = %run_id, status = %final_status, error = %err, "run stopped");
1496
1497 self.publish_run_status_changed(
1498 workflow_name,
1499 run_id,
1500 final_status,
1501 Some(err.to_string()),
1502 ctx,
1503 total_duration,
1504 run_labels,
1505 );
1506
1507 #[cfg(feature = "prometheus")]
1508 self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
1509
1510 return Err(err);
1511 }
1512 }
1513
1514 self.publish_run_status_changed(
1515 workflow_name,
1516 run_id,
1517 final_status,
1518 None,
1519 ctx,
1520 total_duration,
1521 run_labels,
1522 );
1523
1524 #[cfg(feature = "prometheus")]
1525 self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
1526
1527 Ok(WorkflowResult {
1528 run: final_run,
1529 steps: ctx.step_results().to_vec(),
1530 })
1531 }
1532
1533 #[cfg(feature = "prometheus")]
1535 fn emit_run_metrics(
1536 &self,
1537 workflow_name: &str,
1538 status: RunStatus,
1539 duration_ms: u64,
1540 ctx: &WorkflowContext,
1541 ) {
1542 let status_str = status.to_string();
1543 let wf = workflow_name.to_string();
1544
1545 counter!(RUNS_TOTAL, "workflow" => wf.clone(), "status" => status_str.clone()).increment(1);
1546 histogram!(RUN_DURATION_SECONDS, "workflow" => wf.clone(), "status" => status_str)
1547 .record(duration_ms as f64 / 1000.0);
1548 histogram!(RUN_COST_USD, "workflow" => wf.clone()).record(
1549 ctx.total_cost_usd()
1550 .to_string()
1551 .parse::<f64>()
1552 .unwrap_or(0.0),
1553 );
1554 gauge!(RUNS_ACTIVE, "workflow" => wf).decrement(1.0);
1555 }
1556
1557 fn on_run_budget_exceeded(&self, workflow_name: &str, run_id: Uuid, err: &EngineError) {
1563 let EngineError::RunBudgetExceeded {
1564 limit_usd,
1565 spent_usd,
1566 step_budget_usd,
1567 ..
1568 } = err
1569 else {
1570 return;
1571 };
1572
1573 #[cfg(feature = "prometheus")]
1574 counter!(
1575 RUN_BUDGET_EXCEEDED_TOTAL,
1576 "workflow" => workflow_name.to_string(),
1577 "scope" => "run",
1578 )
1579 .increment(1);
1580
1581 self.event_publisher
1582 .publish(Event::RunBudgetExceeded(RunBudgetExceededEvent {
1583 run_id,
1584 workflow_name: workflow_name.to_string(),
1585 limit_usd: *limit_usd,
1586 spent_usd: *spent_usd,
1587 step_budget_usd: *step_budget_usd,
1588 at: Utc::now(),
1589 }));
1590 }
1591
1592 #[allow(clippy::too_many_arguments)]
1597 fn publish_run_status_changed(
1598 &self,
1599 workflow_name: &str,
1600 run_id: Uuid,
1601 to: RunStatus,
1602 error: Option<String>,
1603 ctx: &WorkflowContext,
1604 duration_ms: u64,
1605 labels: HashMap<String, String>,
1606 ) {
1607 let now = Utc::now();
1608 let cost_usd = ctx.total_cost_usd();
1609 let wf = workflow_name.to_string();
1610
1611 self.event_publisher
1612 .publish(Event::RunStatusChanged(RunStatusChangedEvent {
1613 run_id,
1614 workflow_name: wf.clone(),
1615 from: RunStatus::Running,
1616 to,
1617 error: error.clone(),
1618 cost_usd,
1619 duration_ms,
1620 labels: labels.clone(),
1621 at: now,
1622 }));
1623
1624 if to == RunStatus::Failed {
1625 self.event_publisher
1626 .publish(Event::RunFailed(RunFailedEvent {
1627 run_id,
1628 workflow_name: wf,
1629 error,
1630 cost_usd,
1631 duration_ms,
1632 labels,
1633 at: now,
1634 }));
1635 }
1636 }
1637}
1638
1639impl fmt::Debug for Engine {
1640 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1641 f.debug_struct("Engine")
1642 .field("handlers", &self.handlers.keys().collect::<Vec<_>>())
1643 .finish_non_exhaustive()
1644 }
1645}
1646
1647#[cfg(test)]
1648mod tests {
1649 use super::*;
1650 use crate::config::ShellConfig;
1651 use crate::handler::{HandlerFuture, WorkflowHandler};
1652 use ironflow_core::providers::claude::ClaudeCodeProvider;
1653 use ironflow_core::providers::record_replay::RecordReplayProvider;
1654 use ironflow_store::memory::InMemoryStore;
1655 use ironflow_store::models::StepStatus;
1656 use serde_json::json;
1657
1658 struct EchoWorkflow;
1660
1661 impl WorkflowHandler for EchoWorkflow {
1662 fn name(&self) -> &str {
1663 "echo-workflow"
1664 }
1665
1666 fn describe(&self) -> WorkflowInfo {
1667 WorkflowInfo {
1668 description: "A simple workflow that echoes hello".to_string(),
1669 source_code: None,
1670 sub_workflows: Vec::new(),
1671 category: None,
1672 version: self.version().map(str::to_string),
1673 compatible_versions: Vec::new(),
1674 input_schema: None,
1675 default_labels: HashMap::new(),
1676 schedule: self.schedule().cloned(),
1677 default_max_cost_usd: self.default_max_cost_usd(),
1678 }
1679 }
1680
1681 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
1682 Box::pin(async move {
1683 ctx.shell("greet", ShellConfig::new("echo hello")).await?;
1684 Ok(())
1685 })
1686 }
1687 }
1688
1689 struct FailingWorkflow;
1691
1692 impl WorkflowHandler for FailingWorkflow {
1693 fn name(&self) -> &str {
1694 "failing-workflow"
1695 }
1696
1697 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
1698 Box::pin(async move {
1699 ctx.shell("fail", ShellConfig::new("exit 1")).await?;
1700 Ok(())
1701 })
1702 }
1703 }
1704
1705 fn create_test_engine() -> Engine {
1706 let store = Arc::new(InMemoryStore::new());
1707 let inner = ClaudeCodeProvider::new();
1708 let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
1709 inner,
1710 "/tmp/ironflow-fixtures",
1711 ));
1712 Engine::new(store, provider)
1713 }
1714
1715 #[test]
1716 fn engine_new_creates_instance() {
1717 let engine = create_test_engine();
1718 assert_eq!(engine.handler_names().len(), 0);
1719 }
1720
1721 #[test]
1722 fn engine_register_handler() {
1723 let mut engine = create_test_engine();
1724 let result = engine.register(EchoWorkflow);
1725 assert!(result.is_ok());
1726 assert_eq!(engine.handler_names().len(), 1);
1727 assert!(engine.handler_names().contains(&"echo-workflow"));
1728 }
1729
1730 #[test]
1731 fn engine_register_duplicate_returns_error() {
1732 let mut engine = create_test_engine();
1733 engine.register(EchoWorkflow).unwrap();
1734 let result = engine.register(EchoWorkflow);
1735 assert!(result.is_err());
1736 }
1737
1738 #[test]
1739 fn engine_get_handler_found() {
1740 let mut engine = create_test_engine();
1741 engine.register(EchoWorkflow).unwrap();
1742 let handler = engine.get_handler("echo-workflow");
1743 assert!(handler.is_some());
1744 }
1745
1746 #[test]
1747 fn engine_get_handler_not_found() {
1748 let engine = create_test_engine();
1749 let handler = engine.get_handler("nonexistent");
1750 assert!(handler.is_none());
1751 }
1752
1753 #[test]
1754 fn engine_handler_names_lists_all() {
1755 let mut engine = create_test_engine();
1756 engine.register(EchoWorkflow).unwrap();
1757 engine.register(FailingWorkflow).unwrap();
1758 let names = engine.handler_names();
1759 assert_eq!(names.len(), 2);
1760 assert!(names.contains(&"echo-workflow"));
1761 assert!(names.contains(&"failing-workflow"));
1762 }
1763
1764 #[test]
1765 fn engine_handler_info_returns_description() {
1766 let mut engine = create_test_engine();
1767 engine.register(EchoWorkflow).unwrap();
1768 let info = engine.handler_info("echo-workflow");
1769 assert!(info.is_some());
1770 let info = info.unwrap();
1771 assert_eq!(info.description, "A simple workflow that echoes hello");
1772 }
1773
1774 struct CategorizedWorkflow;
1775
1776 impl WorkflowHandler for CategorizedWorkflow {
1777 fn name(&self) -> &str {
1778 "categorized"
1779 }
1780 fn category(&self) -> Option<&str> {
1781 Some("data/etl")
1782 }
1783 fn execute<'a>(
1784 &'a self,
1785 _ctx: &'a mut WorkflowContext,
1786 ) -> crate::handler::HandlerFuture<'a> {
1787 Box::pin(async move { Ok(()) })
1788 }
1789 }
1790
1791 #[test]
1792 fn engine_default_describe_propagates_category() {
1793 let mut engine = create_test_engine();
1794 engine.register(CategorizedWorkflow).unwrap();
1795 let info = engine.handler_info("categorized").unwrap();
1796 assert_eq!(info.category.as_deref(), Some("data/etl"));
1797 }
1798
1799 #[test]
1800 fn engine_default_describe_without_category() {
1801 let mut engine = create_test_engine();
1802 engine.register(EchoWorkflow).unwrap();
1803 let info = engine.handler_info("echo-workflow").unwrap();
1804 assert!(info.category.is_none());
1805 }
1806
1807 struct ScheduledWorkflow {
1812 schedule: CronSchedule,
1813 }
1814
1815 impl ScheduledWorkflow {
1816 fn new() -> Self {
1817 Self {
1818 schedule: CronSchedule::new("0 0 * * * *").unwrap(),
1819 }
1820 }
1821 }
1822
1823 impl WorkflowHandler for ScheduledWorkflow {
1824 fn name(&self) -> &str {
1825 "scheduled"
1826 }
1827 fn schedule(&self) -> Option<&CronSchedule> {
1828 Some(&self.schedule)
1829 }
1830 fn execute<'a>(
1831 &'a self,
1832 _ctx: &'a mut WorkflowContext,
1833 ) -> crate::handler::HandlerFuture<'a> {
1834 Box::pin(async move { Ok(()) })
1835 }
1836 }
1837
1838 #[test]
1839 fn engine_default_describe_propagates_schedule() {
1840 let mut engine = create_test_engine();
1841 engine.register(ScheduledWorkflow::new()).unwrap();
1842 let info = engine.handler_info("scheduled").unwrap();
1843 assert_eq!(
1844 info.schedule.as_ref().map(|s| s.as_str()),
1845 Some("0 0 * * * *")
1846 );
1847 }
1848
1849 #[test]
1850 fn engine_default_describe_without_schedule() {
1851 let mut engine = create_test_engine();
1852 engine.register(EchoWorkflow).unwrap();
1853 let info = engine.handler_info("echo-workflow").unwrap();
1854 assert!(info.schedule.is_none());
1855 }
1856
1857 #[test]
1858 fn scheduled_handlers_returns_only_scheduled() {
1859 let mut engine = create_test_engine();
1860 engine.register(EchoWorkflow).unwrap();
1861 engine.register(ScheduledWorkflow::new()).unwrap();
1862 engine.register(FailingWorkflow).unwrap();
1863
1864 let scheduled = engine.scheduled_handlers();
1865 assert_eq!(scheduled.len(), 1);
1866 assert_eq!(scheduled[0].0, "scheduled");
1867 assert_eq!(scheduled[0].1.as_str(), "0 0 * * * *");
1868 }
1869
1870 #[test]
1871 fn scheduled_handlers_empty_when_none_scheduled() {
1872 let mut engine = create_test_engine();
1873 engine.register(EchoWorkflow).unwrap();
1874 engine.register(FailingWorkflow).unwrap();
1875
1876 let scheduled = engine.scheduled_handlers();
1877 assert!(scheduled.is_empty());
1878 }
1879
1880 struct BadCategoryWorkflow(&'static str);
1881
1882 impl WorkflowHandler for BadCategoryWorkflow {
1883 fn name(&self) -> &str {
1884 "bad-category"
1885 }
1886 fn category(&self) -> Option<&str> {
1887 Some(self.0)
1888 }
1889 fn execute<'a>(
1890 &'a self,
1891 _ctx: &'a mut WorkflowContext,
1892 ) -> crate::handler::HandlerFuture<'a> {
1893 Box::pin(async move { Ok(()) })
1894 }
1895 }
1896
1897 #[test]
1898 fn engine_register_rejects_empty_category() {
1899 let mut engine = create_test_engine();
1900 let err = engine.register(BadCategoryWorkflow("")).unwrap_err();
1901 match err {
1902 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty category")),
1903 other => panic!("expected InvalidWorkflow, got {other:?}"),
1904 }
1905 }
1906
1907 #[test]
1908 fn engine_register_rejects_leading_slash_category() {
1909 let mut engine = create_test_engine();
1910 let err = engine
1911 .register(BadCategoryWorkflow("/data/etl"))
1912 .unwrap_err();
1913 match err {
1914 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("leading '/'")),
1915 other => panic!("expected InvalidWorkflow, got {other:?}"),
1916 }
1917 }
1918
1919 #[test]
1920 fn engine_register_rejects_trailing_slash_category() {
1921 let mut engine = create_test_engine();
1922 let err = engine
1923 .register(BadCategoryWorkflow("data/etl/"))
1924 .unwrap_err();
1925 match err {
1926 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("trailing '/'")),
1927 other => panic!("expected InvalidWorkflow, got {other:?}"),
1928 }
1929 }
1930
1931 #[test]
1932 fn engine_register_rejects_double_slash_category() {
1933 let mut engine = create_test_engine();
1934 let err = engine
1935 .register(BadCategoryWorkflow("data//etl"))
1936 .unwrap_err();
1937 match err {
1938 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty segment")),
1939 other => panic!("expected InvalidWorkflow, got {other:?}"),
1940 }
1941 }
1942
1943 #[test]
1944 fn engine_register_rejects_whitespace_only_segment_category() {
1945 let mut engine = create_test_engine();
1946 let err = engine
1947 .register(BadCategoryWorkflow("data/ /etl"))
1948 .unwrap_err();
1949 match err {
1950 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("whitespace-only segment")),
1951 other => panic!("expected InvalidWorkflow, got {other:?}"),
1952 }
1953 }
1954
1955 #[test]
1956 fn engine_register_accepts_valid_nested_category() {
1957 let mut engine = create_test_engine();
1958 assert!(engine.register(CategorizedWorkflow).is_ok());
1959 }
1960
1961 #[tokio::test]
1962 async fn engine_unknown_workflow_returns_error() {
1963 let engine = create_test_engine();
1964 let result = engine
1965 .run_handler("unknown", TriggerKind::Manual, json!({}))
1966 .await;
1967 assert!(result.is_err());
1968 match result {
1969 Err(EngineError::InvalidWorkflow(msg)) => {
1970 assert!(msg.contains("no handler registered"));
1971 }
1972 _ => panic!("expected InvalidWorkflow error"),
1973 }
1974 }
1975
1976 #[tokio::test]
1977 async fn engine_enqueue_handler_creates_pending_run() {
1978 let mut engine = create_test_engine();
1979 engine.register(EchoWorkflow).unwrap();
1980
1981 let run = engine
1982 .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
1983 .await
1984 .unwrap();
1985 assert_eq!(run.status.state, RunStatus::Pending);
1986 assert_eq!(run.workflow_name, "echo-workflow");
1987 }
1988
1989 #[tokio::test]
1990 async fn enqueue_handler_leaves_the_run_unattributed() {
1991 let mut engine = create_test_engine();
1992 engine.register(EchoWorkflow).unwrap();
1993
1994 let run = engine
1995 .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
1996 .await
1997 .unwrap();
1998
1999 assert!(run.created_by.is_none());
2000 }
2001
2002 #[tokio::test]
2003 async fn enqueue_handler_with_options_records_the_author() {
2004 let mut engine = create_test_engine();
2005 engine.register(EchoWorkflow).unwrap();
2006 let actor = RunActor::User {
2007 user_id: Uuid::now_v7(),
2008 };
2009
2010 let run = engine
2011 .enqueue_handler_with_options(
2012 "echo-workflow",
2013 TriggerKind::Api,
2014 json!({}),
2015 EnqueueOptions {
2016 created_by: Some(actor.clone()),
2017 ..Default::default()
2018 },
2019 )
2020 .await
2021 .unwrap()
2022 .into_run();
2023
2024 assert_eq!(run.created_by, Some(actor));
2025 }
2026
2027 #[tokio::test]
2028 async fn enqueue_handler_with_options_accepts_no_author() {
2029 let mut engine = create_test_engine();
2030 engine.register(EchoWorkflow).unwrap();
2031
2032 let run = engine
2033 .enqueue_handler_with_options(
2034 "echo-workflow",
2035 TriggerKind::Cron {
2036 schedule: "0 * * * * *".to_string(),
2037 },
2038 json!({}),
2039 EnqueueOptions::default(),
2040 )
2041 .await
2042 .unwrap()
2043 .into_run();
2044
2045 assert!(run.created_by.is_none());
2046 }
2047
2048 #[tokio::test]
2049 async fn run_handler_leaves_the_run_unattributed() {
2050 let mut engine = create_test_engine();
2051 engine.register(EchoWorkflow).unwrap();
2052
2053 let run = engine
2054 .run_handler("echo-workflow", TriggerKind::Manual, json!({}))
2055 .await
2056 .unwrap()
2057 .run;
2058
2059 assert!(run.created_by.is_none());
2060 }
2061
2062 #[tokio::test]
2063 async fn engine_register_boxed() {
2064 let mut engine = create_test_engine();
2065 let handler: Box<dyn WorkflowHandler> = Box::new(EchoWorkflow);
2066 let result = engine.register_boxed(handler);
2067 assert!(result.is_ok());
2068 assert_eq!(engine.handler_names().len(), 1);
2069 }
2070
2071 #[tokio::test]
2072 async fn engine_store_and_provider_accessors() {
2073 let store = Arc::new(InMemoryStore::new());
2074 let inner = ClaudeCodeProvider::new();
2075 let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
2076 inner,
2077 "/tmp/ironflow-fixtures",
2078 ));
2079 let engine = Engine::new(store.clone(), provider.clone());
2080
2081 let _ = engine.store();
2083 let _ = engine.provider();
2084 }
2085
2086 use crate::operation::{Operation, OperationContext};
2091 use async_trait::async_trait;
2092 use ironflow_core::error::OperationError;
2093 use ironflow_store::models::StepKind;
2094
2095 struct FakeGitlabOp {
2096 project_id: u64,
2097 title: String,
2098 }
2099
2100 #[async_trait]
2101 impl Operation for FakeGitlabOp {
2102 fn kind(&self) -> &str {
2103 "gitlab"
2104 }
2105
2106 async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
2107 Ok(json!({
2108 "issue_id": 42,
2109 "project_id": self.project_id,
2110 "title": self.title,
2111 }))
2112 }
2113
2114 fn input(&self) -> Option<Value> {
2115 Some(json!({
2116 "project_id": self.project_id,
2117 "title": self.title,
2118 }))
2119 }
2120 }
2121
2122 struct FailingOp;
2123
2124 #[async_trait]
2125 impl Operation for FailingOp {
2126 fn kind(&self) -> &str {
2127 "broken-service"
2128 }
2129
2130 async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
2131 Err(OperationError::Http {
2132 status: None,
2133 message: "service unavailable".to_string(),
2134 })
2135 }
2136 }
2137
2138 struct OperationWorkflow;
2139
2140 impl WorkflowHandler for OperationWorkflow {
2141 fn name(&self) -> &str {
2142 "operation-workflow"
2143 }
2144
2145 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2146 Box::pin(async move {
2147 let op = FakeGitlabOp {
2148 project_id: 123,
2149 title: "Bug report".to_string(),
2150 };
2151 ctx.operation("create-issue", &op).await?;
2152 Ok(())
2153 })
2154 }
2155 }
2156
2157 struct FailingOperationWorkflow;
2158
2159 impl WorkflowHandler for FailingOperationWorkflow {
2160 fn name(&self) -> &str {
2161 "failing-operation-workflow"
2162 }
2163
2164 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2165 Box::pin(async move {
2166 ctx.operation("broken-call", &FailingOp).await?;
2167 Ok(())
2168 })
2169 }
2170 }
2171
2172 struct MixedWorkflow;
2173
2174 impl WorkflowHandler for MixedWorkflow {
2175 fn name(&self) -> &str {
2176 "mixed-workflow"
2177 }
2178
2179 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2180 Box::pin(async move {
2181 ctx.shell("build", ShellConfig::new("echo built")).await?;
2182 let op = FakeGitlabOp {
2183 project_id: 456,
2184 title: "Deploy done".to_string(),
2185 };
2186 let result = ctx.operation("notify-gitlab", &op).await?;
2187 assert_eq!(result.output["issue_id"], 42);
2188 Ok(())
2189 })
2190 }
2191 }
2192
2193 #[tokio::test]
2194 async fn operation_step_happy_path() {
2195 let mut engine = create_test_engine();
2196 engine.register(OperationWorkflow).unwrap();
2197
2198 let run = engine
2199 .run_handler("operation-workflow", TriggerKind::Manual, json!({}))
2200 .await
2201 .unwrap()
2202 .run;
2203
2204 assert_eq!(run.status.state, RunStatus::Completed);
2205
2206 let steps = engine.store().list_steps(run.id).await.unwrap();
2207
2208 assert_eq!(steps.len(), 1);
2209 assert_eq!(steps[0].name, "create-issue");
2210 assert_eq!(steps[0].kind, StepKind::Custom("gitlab".to_string()));
2211 assert_eq!(
2212 steps[0].status.state,
2213 ironflow_store::models::StepStatus::Completed
2214 );
2215
2216 let output = steps[0].output.as_ref().unwrap();
2217 assert_eq!(output["issue_id"], 42);
2218 assert_eq!(output["project_id"], 123);
2219
2220 let input = steps[0].input.as_ref().unwrap();
2221 assert_eq!(input["project_id"], 123);
2222 assert_eq!(input["title"], "Bug report");
2223 }
2224
2225 #[tokio::test]
2226 async fn operation_step_failure_marks_run_failed() {
2227 let mut engine = create_test_engine();
2228 engine.register(FailingOperationWorkflow).unwrap();
2229
2230 let result = engine
2231 .run_handler("failing-operation-workflow", TriggerKind::Manual, json!({}))
2232 .await;
2233
2234 assert!(result.is_err());
2235 }
2236
2237 #[tokio::test]
2238 async fn operation_mixed_with_shell_steps() {
2239 let mut engine = create_test_engine();
2240 engine.register(MixedWorkflow).unwrap();
2241
2242 let run = engine
2243 .run_handler("mixed-workflow", TriggerKind::Manual, json!({}))
2244 .await
2245 .unwrap()
2246 .run;
2247
2248 assert_eq!(run.status.state, RunStatus::Completed);
2249
2250 let steps = engine.store().list_steps(run.id).await.unwrap();
2251
2252 assert_eq!(steps.len(), 2);
2253 assert_eq!(steps[0].kind, StepKind::Shell);
2254 assert_eq!(steps[1].kind, StepKind::Custom("gitlab".to_string()));
2255 assert_eq!(steps[0].position, 0);
2256 assert_eq!(steps[1].position, 1);
2257 }
2258
2259 use crate::config::ApprovalConfig;
2264
2265 struct SingleApprovalWorkflow;
2266
2267 impl WorkflowHandler for SingleApprovalWorkflow {
2268 fn name(&self) -> &str {
2269 "single-approval"
2270 }
2271
2272 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2273 Box::pin(async move {
2274 ctx.shell("build", ShellConfig::new("echo built")).await?;
2275 ctx.approval("gate", ApprovalConfig::new("OK?")).await?;
2276 ctx.shell("deploy", ShellConfig::new("echo deployed"))
2277 .await?;
2278 Ok(())
2279 })
2280 }
2281 }
2282
2283 struct DoubleApprovalWorkflow;
2284
2285 impl WorkflowHandler for DoubleApprovalWorkflow {
2286 fn name(&self) -> &str {
2287 "double-approval"
2288 }
2289
2290 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2291 Box::pin(async move {
2292 ctx.shell("build", ShellConfig::new("echo built")).await?;
2293 ctx.approval("staging-gate", ApprovalConfig::new("Deploy staging?"))
2294 .await?;
2295 ctx.shell("deploy-staging", ShellConfig::new("echo staging"))
2296 .await?;
2297 ctx.approval("prod-gate", ApprovalConfig::new("Deploy prod?"))
2298 .await?;
2299 ctx.shell("deploy-prod", ShellConfig::new("echo prod"))
2300 .await?;
2301 Ok(())
2302 })
2303 }
2304 }
2305
2306 #[tokio::test]
2307 async fn approval_pauses_run() {
2308 let mut engine = create_test_engine();
2309 engine.register(SingleApprovalWorkflow).unwrap();
2310
2311 let run = engine
2312 .run_handler("single-approval", TriggerKind::Manual, json!({}))
2313 .await
2314 .unwrap()
2315 .run;
2316
2317 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
2318
2319 let steps = engine.store().list_steps(run.id).await.unwrap();
2320 assert_eq!(steps.len(), 2); assert_eq!(steps[0].kind, StepKind::Shell);
2322 assert_eq!(steps[0].status.state, StepStatus::Completed);
2323 assert_eq!(steps[1].kind, StepKind::Approval);
2324 assert_eq!(steps[1].status.state, StepStatus::AwaitingApproval);
2325 }
2326
2327 #[tokio::test]
2328 async fn approval_resume_completes_run() {
2329 let mut engine = create_test_engine();
2330 engine.register(SingleApprovalWorkflow).unwrap();
2331
2332 let run = engine
2334 .run_handler("single-approval", TriggerKind::Manual, json!({}))
2335 .await
2336 .unwrap()
2337 .run;
2338 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
2339
2340 engine
2342 .store()
2343 .update_run_status(run.id, RunStatus::Running)
2344 .await
2345 .unwrap();
2346
2347 let resumed = engine.resume_run(run.id).await.unwrap().run;
2349 assert_eq!(resumed.status.state, RunStatus::Completed);
2350
2351 let steps = engine.store().list_steps(run.id).await.unwrap();
2352 assert_eq!(steps.len(), 3); assert_eq!(steps[0].name, "build");
2354 assert_eq!(steps[0].status.state, StepStatus::Completed);
2355 assert_eq!(steps[1].name, "gate");
2356 assert_eq!(steps[1].kind, StepKind::Approval);
2357 assert_eq!(steps[1].status.state, StepStatus::Completed);
2358 assert_eq!(steps[2].name, "deploy");
2359 assert_eq!(steps[2].status.state, StepStatus::Completed);
2360 }
2361
2362 #[tokio::test]
2363 async fn double_approval_two_resumes() {
2364 let mut engine = create_test_engine();
2365 engine.register(DoubleApprovalWorkflow).unwrap();
2366
2367 let run = engine
2369 .run_handler("double-approval", TriggerKind::Manual, json!({}))
2370 .await
2371 .unwrap()
2372 .run;
2373 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
2374
2375 let steps = engine.store().list_steps(run.id).await.unwrap();
2376 assert_eq!(steps.len(), 2); engine
2380 .store()
2381 .update_run_status(run.id, RunStatus::Running)
2382 .await
2383 .unwrap();
2384
2385 let resumed = engine.resume_run(run.id).await.unwrap().run;
2386 assert_eq!(resumed.status.state, RunStatus::AwaitingApproval);
2387
2388 let steps = engine.store().list_steps(run.id).await.unwrap();
2389 assert_eq!(steps.len(), 4); engine
2393 .store()
2394 .update_run_status(run.id, RunStatus::Running)
2395 .await
2396 .unwrap();
2397
2398 let final_run = engine.resume_run(run.id).await.unwrap().run;
2399 assert_eq!(final_run.status.state, RunStatus::Completed);
2400
2401 let steps = engine.store().list_steps(run.id).await.unwrap();
2402 assert_eq!(steps.len(), 5);
2403 assert_eq!(steps[0].name, "build");
2404 assert_eq!(steps[1].name, "staging-gate");
2405 assert_eq!(steps[2].name, "deploy-staging");
2406 assert_eq!(steps[3].name, "prod-gate");
2407 assert_eq!(steps[4].name, "deploy-prod");
2408
2409 for step in &steps {
2410 assert_eq!(step.status.state, StepStatus::Completed);
2411 }
2412 }
2413
2414 use ironflow_store::models::{NewStep, StepUpdate, step_trace_id};
2419
2420 async fn create_step_with_status(
2421 store: &Arc<dyn Store>,
2422 run_id: Uuid,
2423 name: &str,
2424 position: u32,
2425 status: StepStatus,
2426 ) -> ironflow_store::models::Step {
2427 let step = store
2428 .create_step(NewStep {
2429 run_id,
2430 trace_id: step_trace_id(run_id, name, position),
2431 name: name.to_string(),
2432 kind: StepKind::Shell,
2433 position,
2434 input: None,
2435 is_error_handler: false,
2436 })
2437 .await
2438 .unwrap();
2439
2440 match status {
2441 StepStatus::Pending => {}
2442 StepStatus::Running => {
2443 store
2444 .update_step(
2445 step.id,
2446 StepUpdate {
2447 status: Some(StepStatus::Running),
2448 ..StepUpdate::default()
2449 },
2450 )
2451 .await
2452 .unwrap();
2453 }
2454 StepStatus::Completed => {
2455 store
2456 .update_step(
2457 step.id,
2458 StepUpdate {
2459 status: Some(StepStatus::Running),
2460 ..StepUpdate::default()
2461 },
2462 )
2463 .await
2464 .unwrap();
2465 store
2466 .update_step(
2467 step.id,
2468 StepUpdate {
2469 status: Some(StepStatus::Completed),
2470 ..StepUpdate::default()
2471 },
2472 )
2473 .await
2474 .unwrap();
2475 }
2476 StepStatus::AwaitingApproval => {
2477 store
2478 .update_step(
2479 step.id,
2480 StepUpdate {
2481 status: Some(StepStatus::Running),
2482 ..StepUpdate::default()
2483 },
2484 )
2485 .await
2486 .unwrap();
2487 store
2488 .update_step(
2489 step.id,
2490 StepUpdate {
2491 status: Some(StepStatus::AwaitingApproval),
2492 ..StepUpdate::default()
2493 },
2494 )
2495 .await
2496 .unwrap();
2497 }
2498 _ => panic!("unsupported status for test helper: {status}"),
2499 }
2500
2501 store.get_step(step.id).await.unwrap().unwrap()
2502 }
2503
2504 #[tokio::test]
2505 async fn fail_orphaned_steps_marks_running_as_failed() {
2506 let engine = create_test_engine();
2507 let run = engine
2508 .store()
2509 .create_run(NewRun {
2510 created_by: None,
2511 workflow_name: "test".to_string(),
2512 trigger: TriggerKind::Manual,
2513 payload: json!({}),
2514 max_retries: 0,
2515 handler_version: None,
2516 labels: HashMap::new(),
2517 scheduled_at: None,
2518 idempotency_key: None,
2519 max_cost_usd: None,
2520 })
2521 .await
2522 .unwrap()
2523 .into_run();
2524
2525 let step = create_step_with_status(
2526 engine.store(),
2527 run.id,
2528 "running-step",
2529 0,
2530 StepStatus::Running,
2531 )
2532 .await;
2533
2534 engine
2535 .fail_orphaned_steps(run.id, "parent run timed out")
2536 .await
2537 .unwrap();
2538
2539 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
2540 assert_eq!(updated.status.state, StepStatus::Failed);
2541 assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
2542 assert!(updated.completed_at.is_some());
2543 }
2544
2545 #[tokio::test]
2546 async fn fail_orphaned_steps_marks_pending_as_skipped() {
2547 let engine = create_test_engine();
2548 let run = engine
2549 .store()
2550 .create_run(NewRun {
2551 created_by: None,
2552 workflow_name: "test".to_string(),
2553 trigger: TriggerKind::Manual,
2554 payload: json!({}),
2555 max_retries: 0,
2556 handler_version: None,
2557 labels: HashMap::new(),
2558 scheduled_at: None,
2559 idempotency_key: None,
2560 max_cost_usd: None,
2561 })
2562 .await
2563 .unwrap()
2564 .into_run();
2565
2566 let step = create_step_with_status(
2567 engine.store(),
2568 run.id,
2569 "pending-step",
2570 0,
2571 StepStatus::Pending,
2572 )
2573 .await;
2574
2575 engine
2576 .fail_orphaned_steps(run.id, "parent run timed out")
2577 .await
2578 .unwrap();
2579
2580 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
2581 assert_eq!(updated.status.state, StepStatus::Skipped);
2582 assert!(updated.error.is_none());
2583 assert!(updated.completed_at.is_some());
2584 }
2585
2586 #[tokio::test]
2587 async fn fail_orphaned_steps_marks_awaiting_approval_as_failed() {
2588 let engine = create_test_engine();
2589 let run = engine
2590 .store()
2591 .create_run(NewRun {
2592 created_by: None,
2593 workflow_name: "test".to_string(),
2594 trigger: TriggerKind::Manual,
2595 payload: json!({}),
2596 max_retries: 0,
2597 handler_version: None,
2598 labels: HashMap::new(),
2599 scheduled_at: None,
2600 idempotency_key: None,
2601 max_cost_usd: None,
2602 })
2603 .await
2604 .unwrap()
2605 .into_run();
2606
2607 let step = create_step_with_status(
2608 engine.store(),
2609 run.id,
2610 "approval-step",
2611 0,
2612 StepStatus::AwaitingApproval,
2613 )
2614 .await;
2615
2616 engine
2617 .fail_orphaned_steps(run.id, "parent run timed out")
2618 .await
2619 .unwrap();
2620
2621 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
2622 assert_eq!(updated.status.state, StepStatus::Failed);
2623 assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
2624 assert!(updated.completed_at.is_some());
2625 }
2626
2627 #[tokio::test]
2628 async fn fail_orphaned_steps_skips_terminal_steps() {
2629 let engine = create_test_engine();
2630 let run = engine
2631 .store()
2632 .create_run(NewRun {
2633 created_by: None,
2634 workflow_name: "test".to_string(),
2635 trigger: TriggerKind::Manual,
2636 payload: json!({}),
2637 max_retries: 0,
2638 handler_version: None,
2639 labels: HashMap::new(),
2640 scheduled_at: None,
2641 idempotency_key: None,
2642 max_cost_usd: None,
2643 })
2644 .await
2645 .unwrap()
2646 .into_run();
2647
2648 let completed_step =
2649 create_step_with_status(engine.store(), run.id, "done", 0, StepStatus::Completed).await;
2650 let running_step =
2651 create_step_with_status(engine.store(), run.id, "in-flight", 1, StepStatus::Running)
2652 .await;
2653
2654 engine
2655 .fail_orphaned_steps(run.id, "parent run timed out")
2656 .await
2657 .unwrap();
2658
2659 let completed = engine
2660 .store()
2661 .get_step(completed_step.id)
2662 .await
2663 .unwrap()
2664 .unwrap();
2665 assert_eq!(completed.status.state, StepStatus::Completed);
2666
2667 let failed = engine
2668 .store()
2669 .get_step(running_step.id)
2670 .await
2671 .unwrap()
2672 .unwrap();
2673 assert_eq!(failed.status.state, StepStatus::Failed);
2674 }
2675
2676 #[tokio::test]
2677 async fn fail_orphaned_steps_mixed_states() {
2678 let engine = create_test_engine();
2679 let run = engine
2680 .store()
2681 .create_run(NewRun {
2682 created_by: None,
2683 workflow_name: "test".to_string(),
2684 trigger: TriggerKind::Manual,
2685 payload: json!({}),
2686 max_retries: 0,
2687 handler_version: None,
2688 labels: HashMap::new(),
2689 scheduled_at: None,
2690 idempotency_key: None,
2691 max_cost_usd: None,
2692 })
2693 .await
2694 .unwrap()
2695 .into_run();
2696
2697 let s_completed =
2698 create_step_with_status(engine.store(), run.id, "step-1", 0, StepStatus::Completed)
2699 .await;
2700 let s_running =
2701 create_step_with_status(engine.store(), run.id, "step-2", 1, StepStatus::Running).await;
2702 let s_pending =
2703 create_step_with_status(engine.store(), run.id, "step-3", 2, StepStatus::Pending).await;
2704
2705 engine.fail_orphaned_steps(run.id, "timeout").await.unwrap();
2706
2707 let r_completed = engine
2708 .store()
2709 .get_step(s_completed.id)
2710 .await
2711 .unwrap()
2712 .unwrap();
2713 assert_eq!(r_completed.status.state, StepStatus::Completed);
2714
2715 let r_running = engine
2716 .store()
2717 .get_step(s_running.id)
2718 .await
2719 .unwrap()
2720 .unwrap();
2721 assert_eq!(r_running.status.state, StepStatus::Failed);
2722 assert_eq!(r_running.error.as_deref(), Some("timeout"));
2723
2724 let r_pending = engine
2725 .store()
2726 .get_step(s_pending.id)
2727 .await
2728 .unwrap()
2729 .unwrap();
2730 assert_eq!(r_pending.status.state, StepStatus::Skipped);
2731 assert!(r_pending.error.is_none());
2732 }
2733
2734 #[tokio::test]
2735 async fn fail_orphaned_steps_no_steps_is_noop() {
2736 let engine = create_test_engine();
2737 let run = engine
2738 .store()
2739 .create_run(NewRun {
2740 created_by: None,
2741 workflow_name: "test".to_string(),
2742 trigger: TriggerKind::Manual,
2743 payload: json!({}),
2744 max_retries: 0,
2745 handler_version: None,
2746 labels: HashMap::new(),
2747 scheduled_at: None,
2748 idempotency_key: None,
2749 max_cost_usd: None,
2750 })
2751 .await
2752 .unwrap()
2753 .into_run();
2754
2755 let result = engine.fail_orphaned_steps(run.id, "timeout").await;
2756 assert!(result.is_ok());
2757 }
2758
2759 #[tokio::test]
2760 async fn fail_orphaned_steps_preserves_existing_error() {
2761 let engine = create_test_engine();
2762 let run = engine
2763 .store()
2764 .create_run(NewRun {
2765 created_by: None,
2766 workflow_name: "test".to_string(),
2767 trigger: TriggerKind::Manual,
2768 payload: json!({}),
2769 max_retries: 0,
2770 handler_version: None,
2771 labels: HashMap::new(),
2772 scheduled_at: None,
2773 idempotency_key: None,
2774 max_cost_usd: None,
2775 })
2776 .await
2777 .unwrap()
2778 .into_run();
2779
2780 let step_with_error = create_step_with_status(
2781 engine.store(),
2782 run.id,
2783 "already-errored",
2784 0,
2785 StepStatus::Running,
2786 )
2787 .await;
2788
2789 engine
2790 .store()
2791 .update_step(
2792 step_with_error.id,
2793 StepUpdate {
2794 error: Some("real error from provider".to_string()),
2795 ..StepUpdate::default()
2796 },
2797 )
2798 .await
2799 .unwrap();
2800
2801 let step_no_error = create_step_with_status(
2802 engine.store(),
2803 run.id,
2804 "no-error-yet",
2805 1,
2806 StepStatus::Running,
2807 )
2808 .await;
2809
2810 engine
2811 .fail_orphaned_steps(run.id, "parent run failed")
2812 .await
2813 .unwrap();
2814
2815 let updated_with = engine
2816 .store()
2817 .get_step(step_with_error.id)
2818 .await
2819 .unwrap()
2820 .unwrap();
2821 assert_eq!(updated_with.status.state, StepStatus::Failed);
2822 assert_eq!(
2823 updated_with.error.as_deref(),
2824 Some("real error from provider"),
2825 );
2826
2827 let updated_without = engine
2828 .store()
2829 .get_step(step_no_error.id)
2830 .await
2831 .unwrap()
2832 .unwrap();
2833 assert_eq!(updated_without.status.state, StepStatus::Failed);
2834 assert_eq!(updated_without.error.as_deref(), Some("parent run failed"),);
2835 }
2836}