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