1use std::collections::HashMap;
10use std::fmt;
11use std::sync::Arc;
12use std::time::Instant;
13
14use chrono::{DateTime, Utc};
15use serde_json::Value;
16use tracing::{error, info, warn};
17use uuid::Uuid;
18
19#[cfg(feature = "prometheus")]
20use ironflow_core::metric_names::{RUN_COST_USD, RUN_DURATION_SECONDS, RUNS_ACTIVE, RUNS_TOTAL};
21use ironflow_core::provider::AgentProvider;
22use ironflow_store::error::StoreError;
23use ironflow_store::models::{
24 NewRun, Run, RunStatus, RunUpdate, StepStatus, StepUpdate, TriggerKind,
25};
26use ironflow_store::store::Store;
27#[cfg(feature = "prometheus")]
28use metrics::{counter, gauge, histogram};
29
30use crate::context::WorkflowContext;
31use crate::error::EngineError;
32use crate::handler::{WorkflowHandler, WorkflowInfo};
33use crate::log_sender::LogSender;
34use crate::notify::{Event, EventPublisher, EventSubscriber};
35use crate::schedule::CronSchedule;
36
37pub struct Engine {
78 store: Arc<dyn Store>,
79 provider: Arc<dyn AgentProvider>,
80 handlers: HashMap<String, Arc<dyn WorkflowHandler>>,
81 event_publisher: EventPublisher,
82 log_sender: Option<LogSender>,
83}
84
85fn validate_category(handler_name: &str, category: &str) -> Result<(), EngineError> {
95 let reject = |reason: &str| {
96 Err(EngineError::InvalidWorkflow(format!(
97 "handler '{handler_name}' has invalid category '{category}': {reason}"
98 )))
99 };
100
101 if category.is_empty() {
102 return reject("empty category");
103 }
104 if category.starts_with('/') {
105 return reject("leading '/'");
106 }
107 if category.ends_with('/') {
108 return reject("trailing '/'");
109 }
110 for segment in category.split('/') {
111 if segment.is_empty() {
112 return reject("empty segment (double '/')");
113 }
114 if segment.trim().is_empty() {
115 return reject("whitespace-only segment");
116 }
117 }
118 Ok(())
119}
120
121impl Engine {
122 pub fn new(store: Arc<dyn Store>, provider: Arc<dyn AgentProvider>) -> Self {
138 Self {
139 store,
140 provider,
141 handlers: HashMap::new(),
142 event_publisher: EventPublisher::new(),
143 log_sender: None,
144 }
145 }
146
147 pub fn set_log_sender(&mut self, sender: LogSender) {
153 self.log_sender = Some(sender);
154 }
155
156 pub fn store(&self) -> &Arc<dyn Store> {
158 &self.store
159 }
160
161 pub fn provider(&self) -> &Arc<dyn AgentProvider> {
163 &self.provider
164 }
165
166 fn build_context(&self, run_id: Uuid) -> WorkflowContext {
168 let handlers = self.handlers.clone();
169 let resolver: crate::context::HandlerResolver =
170 Arc::new(move |name: &str| handlers.get(name).cloned());
171 let mut ctx = WorkflowContext::with_handler_resolver(
172 run_id,
173 self.store.clone(),
174 self.provider.clone(),
175 resolver,
176 );
177 if let Some(ref sender) = self.log_sender {
178 ctx.set_log_sender(sender.clone());
179 }
180 ctx
181 }
182
183 pub fn register(&mut self, handler: impl WorkflowHandler + 'static) -> Result<(), EngineError> {
227 let name = handler.name().to_string();
228 if self.handlers.contains_key(&name) {
229 return Err(EngineError::InvalidWorkflow(format!(
230 "handler '{}' already registered",
231 name
232 )));
233 }
234 if let Some(category) = handler.category() {
235 validate_category(&name, category)?;
236 }
237 self.handlers.insert(name, Arc::new(handler));
238 Ok(())
239 }
240
241 pub fn register_boxed(&mut self, handler: Box<dyn WorkflowHandler>) -> Result<(), EngineError> {
248 let name = handler.name().to_string();
249 if self.handlers.contains_key(&name) {
250 return Err(EngineError::InvalidWorkflow(format!(
251 "handler '{}' already registered",
252 name
253 )));
254 }
255 if let Some(category) = handler.category() {
256 validate_category(&name, category)?;
257 }
258 self.handlers.insert(name, Arc::from(handler));
259 Ok(())
260 }
261
262 pub fn get_handler(&self, name: &str) -> Option<&Arc<dyn WorkflowHandler>> {
264 self.handlers.get(name)
265 }
266
267 pub fn handler_names(&self) -> Vec<&str> {
269 self.handlers.keys().map(|s| s.as_str()).collect()
270 }
271
272 pub fn handler_info(&self, name: &str) -> Option<WorkflowInfo> {
274 self.handlers.get(name).map(|h| h.describe())
275 }
276
277 pub fn scheduled_handlers(&self) -> Vec<(&str, &CronSchedule)> {
301 self.handlers
302 .iter()
303 .filter_map(|(name, handler)| handler.schedule().map(|sched| (name.as_str(), sched)))
304 .collect()
305 }
306
307 pub fn subscribe(
332 &mut self,
333 subscriber: impl EventSubscriber + 'static,
334 event_types: &[&'static str],
335 ) {
336 self.event_publisher.subscribe(subscriber, event_types);
337 }
338
339 pub fn event_publisher(&self) -> &EventPublisher {
344 &self.event_publisher
345 }
346
347 #[tracing::instrument(name = "engine.run_handler", skip_all, fields(workflow = %handler_name))]
377 pub async fn run_handler(
378 &self,
379 handler_name: &str,
380 trigger: TriggerKind,
381 payload: Value,
382 ) -> Result<Run, EngineError> {
383 let handler = self
384 .handlers
385 .get(handler_name)
386 .ok_or_else(|| {
387 EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
388 })?
389 .clone();
390
391 let handler_version = handler.version().map(str::to_string);
392 let run = self
393 .store
394 .create_run(NewRun {
395 workflow_name: handler_name.to_string(),
396 trigger,
397 payload,
398 max_retries: 0,
399 handler_version,
400 labels: handler.default_labels(),
401 scheduled_at: None,
402 })
403 .await?;
404
405 let run_id = run.id;
406 info!(run_id = %run_id, handler_version = run.handler_version.as_deref().unwrap_or(""), "run created");
407
408 self.store
409 .update_run_status(run_id, RunStatus::Running)
410 .await?;
411
412 #[cfg(feature = "prometheus")]
413 gauge!(RUNS_ACTIVE, "workflow" => handler_name.to_string()).increment(1.0);
414
415 let run_start = Instant::now();
416 let mut ctx = self.build_context(run_id);
417
418 let result = handler.execute(&mut ctx).await;
419 self.finalize_run(run_id, handler_name, result, &ctx, run_start)
420 .await
421 }
422
423 #[tracing::instrument(name = "engine.enqueue_handler", skip_all, fields(workflow = %handler_name))]
432 pub async fn enqueue_handler(
433 &self,
434 handler_name: &str,
435 trigger: TriggerKind,
436 payload: Value,
437 max_retries: u32,
438 ) -> Result<Run, EngineError> {
439 self.enqueue_handler_with_options(
440 handler_name,
441 trigger,
442 payload,
443 max_retries,
444 HashMap::new(),
445 None,
446 )
447 .await
448 }
449
450 #[tracing::instrument(name = "engine.enqueue_handler_with_options", skip_all, fields(workflow = %handler_name))]
456 pub async fn enqueue_handler_with_options(
457 &self,
458 handler_name: &str,
459 trigger: TriggerKind,
460 payload: Value,
461 max_retries: u32,
462 labels: HashMap<String, String>,
463 scheduled_at: Option<DateTime<Utc>>,
464 ) -> Result<Run, EngineError> {
465 let handler = self.handlers.get(handler_name).ok_or_else(|| {
466 EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
467 })?;
468
469 let handler_version = handler.version().map(str::to_string);
470 let mut merged_labels = handler.default_labels();
471 merged_labels.extend(labels);
472
473 let run = self
474 .store
475 .create_run(NewRun {
476 workflow_name: handler_name.to_string(),
477 trigger,
478 payload,
479 max_retries,
480 handler_version,
481 labels: merged_labels,
482 scheduled_at,
483 })
484 .await?;
485
486 info!(run_id = %run.id, workflow = %handler_name, "handler run enqueued");
487 Ok(run)
488 }
489
490 #[tracing::instrument(name = "engine.execute_handler_run", skip_all, fields(run_id = %run_id))]
499 pub async fn execute_handler_run(&self, run_id: Uuid) -> Result<Run, EngineError> {
500 let run = self
501 .store
502 .get_run(run_id)
503 .await?
504 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
505
506 let handler = self
507 .handlers
508 .get(&run.workflow_name)
509 .ok_or_else(|| {
510 EngineError::InvalidWorkflow(format!(
511 "no handler registered: {}",
512 run.workflow_name
513 ))
514 })?
515 .clone();
516
517 #[cfg(feature = "prometheus")]
518 gauge!(RUNS_ACTIVE, "workflow" => run.workflow_name.clone()).increment(1.0);
519
520 let run_start = Instant::now();
521 let mut ctx = self.build_context(run_id);
522
523 let result = handler.execute(&mut ctx).await;
524 self.finalize_run(run_id, &run.workflow_name, result, &ctx, run_start)
525 .await
526 }
527
528 #[tracing::instrument(name = "engine.execute_run", skip_all, fields(run_id = %run_id))]
536 pub async fn execute_run(&self, run_id: Uuid) -> Result<Run, EngineError> {
537 self.execute_handler_run(run_id).await
538 }
539
540 #[tracing::instrument(name = "engine.resume_run", skip_all, fields(run_id = %run_id))]
554 pub async fn resume_run(&self, run_id: Uuid) -> Result<Run, EngineError> {
555 let run = self
556 .store
557 .get_run(run_id)
558 .await?
559 .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
560
561 let handler = self
562 .handlers
563 .get(&run.workflow_name)
564 .ok_or_else(|| {
565 EngineError::InvalidWorkflow(format!(
566 "no handler registered: {}",
567 run.workflow_name
568 ))
569 })?
570 .clone();
571
572 info!(run_id = %run_id, workflow = %run.workflow_name, "resuming run after approval");
573
574 let run_start = Instant::now();
575 let mut ctx = self.build_context(run_id);
576 ctx.load_replay_steps().await?;
577
578 let result = handler.execute(&mut ctx).await;
579 self.finalize_run(run_id, &run.workflow_name, result, &ctx, run_start)
580 .await
581 }
582
583 pub async fn fail_orphaned_steps(
597 &self,
598 run_id: Uuid,
599 error_message: &str,
600 ) -> Result<(), EngineError> {
601 let steps = self.store.list_steps(run_id).await?;
602 let now = Utc::now();
603
604 for step in steps {
605 if step.status.state.is_terminal() {
606 continue;
607 }
608
609 let (target_status, error) = match step.status.state {
610 StepStatus::Running | StepStatus::AwaitingApproval => {
611 let err = if step.error.is_some() {
612 None
613 } else {
614 Some(error_message.to_string())
615 };
616 (StepStatus::Failed, err)
617 }
618 StepStatus::Pending => (StepStatus::Skipped, None),
619 _ => continue,
620 };
621
622 if let Err(e) = self
623 .store
624 .update_step(
625 step.id,
626 StepUpdate {
627 status: Some(target_status),
628 error,
629 completed_at: Some(now),
630 ..StepUpdate::default()
631 },
632 )
633 .await
634 {
635 warn!(
636 run_id = %run_id,
637 step_id = %step.id,
638 step_name = %step.name,
639 error = %e,
640 "failed to cleanup orphaned step"
641 );
642 } else {
643 info!(
644 run_id = %run_id,
645 step_id = %step.id,
646 step_name = %step.name,
647 from = %step.status.state,
648 to = %target_status,
649 "cleaned up orphaned step"
650 );
651 }
652 }
653
654 Ok(())
655 }
656
657 async fn finalize_run(
663 &self,
664 run_id: Uuid,
665 workflow_name: &str,
666 result: Result<(), EngineError>,
667 ctx: &WorkflowContext,
668 run_start: Instant,
669 ) -> Result<Run, EngineError> {
670 let total_duration = run_start.elapsed().as_millis() as u64;
671 let completed_at = Utc::now();
672
673 let final_status;
674 let final_run;
675
676 match result {
677 Ok(()) => {
678 final_status = RunStatus::Completed;
679 final_run = self
680 .store
681 .update_run_returning(
682 run_id,
683 RunUpdate {
684 status: Some(RunStatus::Completed),
685 cost_usd: Some(ctx.total_cost_usd()),
686 duration_ms: Some(total_duration),
687 completed_at: Some(completed_at),
688 ..RunUpdate::default()
689 },
690 )
691 .await?;
692
693 info!(
694 run_id = %run_id,
695 cost_usd = %ctx.total_cost_usd(),
696 duration_ms = total_duration,
697 "run completed"
698 );
699 }
700 Err(EngineError::ApprovalRequired {
701 run_id: approval_run_id,
702 step_id,
703 ref message,
704 }) => {
705 final_status = RunStatus::AwaitingApproval;
706 final_run = self
707 .store
708 .update_run_returning(
709 run_id,
710 RunUpdate {
711 status: Some(RunStatus::AwaitingApproval),
712 cost_usd: Some(ctx.total_cost_usd()),
713 duration_ms: Some(total_duration),
714 ..RunUpdate::default()
715 },
716 )
717 .await?;
718
719 info!(
720 run_id = %approval_run_id,
721 step_id = %step_id,
722 message = %message,
723 "run awaiting approval"
724 );
725 }
726 Err(err) => {
727 final_status = RunStatus::Failed;
728 if let Err(store_err) = self
729 .store
730 .update_run(
731 run_id,
732 RunUpdate {
733 status: Some(RunStatus::Failed),
734 error: Some(err.to_string()),
735 cost_usd: Some(ctx.total_cost_usd()),
736 duration_ms: Some(total_duration),
737 completed_at: Some(completed_at),
738 ..RunUpdate::default()
739 },
740 )
741 .await
742 {
743 error!(run_id = %run_id, store_error = %store_err, "failed to persist run failure");
744 }
745
746 error!(run_id = %run_id, error = %err, "run failed");
747
748 self.publish_run_status_changed(
749 workflow_name,
750 run_id,
751 final_status,
752 Some(err.to_string()),
753 ctx,
754 total_duration,
755 );
756
757 #[cfg(feature = "prometheus")]
758 self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
759
760 return Err(err);
761 }
762 }
763
764 self.publish_run_status_changed(
765 workflow_name,
766 run_id,
767 final_status,
768 None,
769 ctx,
770 total_duration,
771 );
772
773 #[cfg(feature = "prometheus")]
774 self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
775
776 Ok(final_run)
777 }
778
779 #[cfg(feature = "prometheus")]
781 fn emit_run_metrics(
782 &self,
783 workflow_name: &str,
784 status: RunStatus,
785 duration_ms: u64,
786 ctx: &WorkflowContext,
787 ) {
788 let status_str = status.to_string();
789 let wf = workflow_name.to_string();
790
791 counter!(RUNS_TOTAL, "workflow" => wf.clone(), "status" => status_str.clone()).increment(1);
792 histogram!(RUN_DURATION_SECONDS, "workflow" => wf.clone(), "status" => status_str)
793 .record(duration_ms as f64 / 1000.0);
794 histogram!(RUN_COST_USD, "workflow" => wf.clone()).record(
795 ctx.total_cost_usd()
796 .to_string()
797 .parse::<f64>()
798 .unwrap_or(0.0),
799 );
800 gauge!(RUNS_ACTIVE, "workflow" => wf).decrement(1.0);
801 }
802
803 fn publish_run_status_changed(
808 &self,
809 workflow_name: &str,
810 run_id: Uuid,
811 to: RunStatus,
812 error: Option<String>,
813 ctx: &WorkflowContext,
814 duration_ms: u64,
815 ) {
816 let now = Utc::now();
817 let cost_usd = ctx.total_cost_usd();
818 let wf = workflow_name.to_string();
819
820 self.event_publisher.publish(Event::RunStatusChanged {
821 run_id,
822 workflow_name: wf.clone(),
823 from: RunStatus::Running,
824 to,
825 error: error.clone(),
826 cost_usd,
827 duration_ms,
828 at: now,
829 });
830
831 if to == RunStatus::Failed {
832 self.event_publisher.publish(Event::RunFailed {
833 run_id,
834 workflow_name: wf,
835 error,
836 cost_usd,
837 duration_ms,
838 at: now,
839 });
840 }
841 }
842}
843
844impl fmt::Debug for Engine {
845 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
846 f.debug_struct("Engine")
847 .field("handlers", &self.handlers.keys().collect::<Vec<_>>())
848 .finish_non_exhaustive()
849 }
850}
851
852#[cfg(test)]
853mod tests {
854 use super::*;
855 use crate::config::ShellConfig;
856 use crate::handler::{HandlerFuture, WorkflowHandler};
857 use ironflow_core::providers::claude::ClaudeCodeProvider;
858 use ironflow_core::providers::record_replay::RecordReplayProvider;
859 use ironflow_store::memory::InMemoryStore;
860 use ironflow_store::models::StepStatus;
861 use serde_json::json;
862
863 struct EchoWorkflow;
865
866 impl WorkflowHandler for EchoWorkflow {
867 fn name(&self) -> &str {
868 "echo-workflow"
869 }
870
871 fn describe(&self) -> WorkflowInfo {
872 WorkflowInfo {
873 description: "A simple workflow that echoes hello".to_string(),
874 source_code: None,
875 sub_workflows: Vec::new(),
876 category: None,
877 version: self.version().map(str::to_string),
878 input_schema: None,
879 default_labels: HashMap::new(),
880 schedule: self.schedule().cloned(),
881 }
882 }
883
884 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
885 Box::pin(async move {
886 ctx.shell("greet", ShellConfig::new("echo hello")).await?;
887 Ok(())
888 })
889 }
890 }
891
892 struct FailingWorkflow;
894
895 impl WorkflowHandler for FailingWorkflow {
896 fn name(&self) -> &str {
897 "failing-workflow"
898 }
899
900 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
901 Box::pin(async move {
902 ctx.shell("fail", ShellConfig::new("exit 1")).await?;
903 Ok(())
904 })
905 }
906 }
907
908 fn create_test_engine() -> Engine {
909 let store = Arc::new(InMemoryStore::new());
910 let inner = ClaudeCodeProvider::new();
911 let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
912 inner,
913 "/tmp/ironflow-fixtures",
914 ));
915 Engine::new(store, provider)
916 }
917
918 #[test]
919 fn engine_new_creates_instance() {
920 let engine = create_test_engine();
921 assert_eq!(engine.handler_names().len(), 0);
922 }
923
924 #[test]
925 fn engine_register_handler() {
926 let mut engine = create_test_engine();
927 let result = engine.register(EchoWorkflow);
928 assert!(result.is_ok());
929 assert_eq!(engine.handler_names().len(), 1);
930 assert!(engine.handler_names().contains(&"echo-workflow"));
931 }
932
933 #[test]
934 fn engine_register_duplicate_returns_error() {
935 let mut engine = create_test_engine();
936 engine.register(EchoWorkflow).unwrap();
937 let result = engine.register(EchoWorkflow);
938 assert!(result.is_err());
939 }
940
941 #[test]
942 fn engine_get_handler_found() {
943 let mut engine = create_test_engine();
944 engine.register(EchoWorkflow).unwrap();
945 let handler = engine.get_handler("echo-workflow");
946 assert!(handler.is_some());
947 }
948
949 #[test]
950 fn engine_get_handler_not_found() {
951 let engine = create_test_engine();
952 let handler = engine.get_handler("nonexistent");
953 assert!(handler.is_none());
954 }
955
956 #[test]
957 fn engine_handler_names_lists_all() {
958 let mut engine = create_test_engine();
959 engine.register(EchoWorkflow).unwrap();
960 engine.register(FailingWorkflow).unwrap();
961 let names = engine.handler_names();
962 assert_eq!(names.len(), 2);
963 assert!(names.contains(&"echo-workflow"));
964 assert!(names.contains(&"failing-workflow"));
965 }
966
967 #[test]
968 fn engine_handler_info_returns_description() {
969 let mut engine = create_test_engine();
970 engine.register(EchoWorkflow).unwrap();
971 let info = engine.handler_info("echo-workflow");
972 assert!(info.is_some());
973 let info = info.unwrap();
974 assert_eq!(info.description, "A simple workflow that echoes hello");
975 }
976
977 struct CategorizedWorkflow;
978
979 impl WorkflowHandler for CategorizedWorkflow {
980 fn name(&self) -> &str {
981 "categorized"
982 }
983 fn category(&self) -> Option<&str> {
984 Some("data/etl")
985 }
986 fn execute<'a>(
987 &'a self,
988 _ctx: &'a mut WorkflowContext,
989 ) -> crate::handler::HandlerFuture<'a> {
990 Box::pin(async move { Ok(()) })
991 }
992 }
993
994 #[test]
995 fn engine_default_describe_propagates_category() {
996 let mut engine = create_test_engine();
997 engine.register(CategorizedWorkflow).unwrap();
998 let info = engine.handler_info("categorized").unwrap();
999 assert_eq!(info.category.as_deref(), Some("data/etl"));
1000 }
1001
1002 #[test]
1003 fn engine_default_describe_without_category() {
1004 let mut engine = create_test_engine();
1005 engine.register(EchoWorkflow).unwrap();
1006 let info = engine.handler_info("echo-workflow").unwrap();
1007 assert!(info.category.is_none());
1008 }
1009
1010 struct ScheduledWorkflow {
1015 schedule: CronSchedule,
1016 }
1017
1018 impl ScheduledWorkflow {
1019 fn new() -> Self {
1020 Self {
1021 schedule: CronSchedule::new("0 0 * * * *").unwrap(),
1022 }
1023 }
1024 }
1025
1026 impl WorkflowHandler for ScheduledWorkflow {
1027 fn name(&self) -> &str {
1028 "scheduled"
1029 }
1030 fn schedule(&self) -> Option<&CronSchedule> {
1031 Some(&self.schedule)
1032 }
1033 fn execute<'a>(
1034 &'a self,
1035 _ctx: &'a mut WorkflowContext,
1036 ) -> crate::handler::HandlerFuture<'a> {
1037 Box::pin(async move { Ok(()) })
1038 }
1039 }
1040
1041 #[test]
1042 fn engine_default_describe_propagates_schedule() {
1043 let mut engine = create_test_engine();
1044 engine.register(ScheduledWorkflow::new()).unwrap();
1045 let info = engine.handler_info("scheduled").unwrap();
1046 assert_eq!(
1047 info.schedule.as_ref().map(|s| s.as_str()),
1048 Some("0 0 * * * *")
1049 );
1050 }
1051
1052 #[test]
1053 fn engine_default_describe_without_schedule() {
1054 let mut engine = create_test_engine();
1055 engine.register(EchoWorkflow).unwrap();
1056 let info = engine.handler_info("echo-workflow").unwrap();
1057 assert!(info.schedule.is_none());
1058 }
1059
1060 #[test]
1061 fn scheduled_handlers_returns_only_scheduled() {
1062 let mut engine = create_test_engine();
1063 engine.register(EchoWorkflow).unwrap();
1064 engine.register(ScheduledWorkflow::new()).unwrap();
1065 engine.register(FailingWorkflow).unwrap();
1066
1067 let scheduled = engine.scheduled_handlers();
1068 assert_eq!(scheduled.len(), 1);
1069 assert_eq!(scheduled[0].0, "scheduled");
1070 assert_eq!(scheduled[0].1.as_str(), "0 0 * * * *");
1071 }
1072
1073 #[test]
1074 fn scheduled_handlers_empty_when_none_scheduled() {
1075 let mut engine = create_test_engine();
1076 engine.register(EchoWorkflow).unwrap();
1077 engine.register(FailingWorkflow).unwrap();
1078
1079 let scheduled = engine.scheduled_handlers();
1080 assert!(scheduled.is_empty());
1081 }
1082
1083 struct BadCategoryWorkflow(&'static str);
1084
1085 impl WorkflowHandler for BadCategoryWorkflow {
1086 fn name(&self) -> &str {
1087 "bad-category"
1088 }
1089 fn category(&self) -> Option<&str> {
1090 Some(self.0)
1091 }
1092 fn execute<'a>(
1093 &'a self,
1094 _ctx: &'a mut WorkflowContext,
1095 ) -> crate::handler::HandlerFuture<'a> {
1096 Box::pin(async move { Ok(()) })
1097 }
1098 }
1099
1100 #[test]
1101 fn engine_register_rejects_empty_category() {
1102 let mut engine = create_test_engine();
1103 let err = engine.register(BadCategoryWorkflow("")).unwrap_err();
1104 match err {
1105 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty category")),
1106 other => panic!("expected InvalidWorkflow, got {other:?}"),
1107 }
1108 }
1109
1110 #[test]
1111 fn engine_register_rejects_leading_slash_category() {
1112 let mut engine = create_test_engine();
1113 let err = engine
1114 .register(BadCategoryWorkflow("/data/etl"))
1115 .unwrap_err();
1116 match err {
1117 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("leading '/'")),
1118 other => panic!("expected InvalidWorkflow, got {other:?}"),
1119 }
1120 }
1121
1122 #[test]
1123 fn engine_register_rejects_trailing_slash_category() {
1124 let mut engine = create_test_engine();
1125 let err = engine
1126 .register(BadCategoryWorkflow("data/etl/"))
1127 .unwrap_err();
1128 match err {
1129 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("trailing '/'")),
1130 other => panic!("expected InvalidWorkflow, got {other:?}"),
1131 }
1132 }
1133
1134 #[test]
1135 fn engine_register_rejects_double_slash_category() {
1136 let mut engine = create_test_engine();
1137 let err = engine
1138 .register(BadCategoryWorkflow("data//etl"))
1139 .unwrap_err();
1140 match err {
1141 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty segment")),
1142 other => panic!("expected InvalidWorkflow, got {other:?}"),
1143 }
1144 }
1145
1146 #[test]
1147 fn engine_register_rejects_whitespace_only_segment_category() {
1148 let mut engine = create_test_engine();
1149 let err = engine
1150 .register(BadCategoryWorkflow("data/ /etl"))
1151 .unwrap_err();
1152 match err {
1153 EngineError::InvalidWorkflow(msg) => assert!(msg.contains("whitespace-only segment")),
1154 other => panic!("expected InvalidWorkflow, got {other:?}"),
1155 }
1156 }
1157
1158 #[test]
1159 fn engine_register_accepts_valid_nested_category() {
1160 let mut engine = create_test_engine();
1161 assert!(engine.register(CategorizedWorkflow).is_ok());
1162 }
1163
1164 #[tokio::test]
1165 async fn engine_unknown_workflow_returns_error() {
1166 let engine = create_test_engine();
1167 let result = engine
1168 .run_handler("unknown", TriggerKind::Manual, json!({}))
1169 .await;
1170 assert!(result.is_err());
1171 match result {
1172 Err(EngineError::InvalidWorkflow(msg)) => {
1173 assert!(msg.contains("no handler registered"));
1174 }
1175 _ => panic!("expected InvalidWorkflow error"),
1176 }
1177 }
1178
1179 #[tokio::test]
1180 async fn engine_enqueue_handler_creates_pending_run() {
1181 let mut engine = create_test_engine();
1182 engine.register(EchoWorkflow).unwrap();
1183
1184 let run = engine
1185 .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
1186 .await
1187 .unwrap();
1188 assert_eq!(run.status.state, RunStatus::Pending);
1189 assert_eq!(run.workflow_name, "echo-workflow");
1190 }
1191
1192 #[tokio::test]
1193 async fn engine_register_boxed() {
1194 let mut engine = create_test_engine();
1195 let handler: Box<dyn WorkflowHandler> = Box::new(EchoWorkflow);
1196 let result = engine.register_boxed(handler);
1197 assert!(result.is_ok());
1198 assert_eq!(engine.handler_names().len(), 1);
1199 }
1200
1201 #[tokio::test]
1202 async fn engine_store_and_provider_accessors() {
1203 let store = Arc::new(InMemoryStore::new());
1204 let inner = ClaudeCodeProvider::new();
1205 let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
1206 inner,
1207 "/tmp/ironflow-fixtures",
1208 ));
1209 let engine = Engine::new(store.clone(), provider.clone());
1210
1211 let _ = engine.store();
1213 let _ = engine.provider();
1214 }
1215
1216 use crate::operation::Operation;
1221 use ironflow_store::models::StepKind;
1222 use std::future::Future;
1223 use std::pin::Pin;
1224
1225 struct FakeGitlabOp {
1226 project_id: u64,
1227 title: String,
1228 }
1229
1230 impl Operation for FakeGitlabOp {
1231 fn kind(&self) -> &str {
1232 "gitlab"
1233 }
1234
1235 fn execute(&self) -> Pin<Box<dyn Future<Output = Result<Value, EngineError>> + Send + '_>> {
1236 Box::pin(async move {
1237 Ok(json!({
1238 "issue_id": 42,
1239 "project_id": self.project_id,
1240 "title": self.title,
1241 }))
1242 })
1243 }
1244
1245 fn input(&self) -> Option<Value> {
1246 Some(json!({
1247 "project_id": self.project_id,
1248 "title": self.title,
1249 }))
1250 }
1251 }
1252
1253 struct FailingOp;
1254
1255 impl Operation for FailingOp {
1256 fn kind(&self) -> &str {
1257 "broken-service"
1258 }
1259
1260 fn execute(&self) -> Pin<Box<dyn Future<Output = Result<Value, EngineError>> + Send + '_>> {
1261 Box::pin(async move { Err(EngineError::StepConfig("service unavailable".to_string())) })
1262 }
1263 }
1264
1265 struct OperationWorkflow;
1266
1267 impl WorkflowHandler for OperationWorkflow {
1268 fn name(&self) -> &str {
1269 "operation-workflow"
1270 }
1271
1272 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
1273 Box::pin(async move {
1274 let op = FakeGitlabOp {
1275 project_id: 123,
1276 title: "Bug report".to_string(),
1277 };
1278 ctx.operation("create-issue", &op).await?;
1279 Ok(())
1280 })
1281 }
1282 }
1283
1284 struct FailingOperationWorkflow;
1285
1286 impl WorkflowHandler for FailingOperationWorkflow {
1287 fn name(&self) -> &str {
1288 "failing-operation-workflow"
1289 }
1290
1291 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
1292 Box::pin(async move {
1293 ctx.operation("broken-call", &FailingOp).await?;
1294 Ok(())
1295 })
1296 }
1297 }
1298
1299 struct MixedWorkflow;
1300
1301 impl WorkflowHandler for MixedWorkflow {
1302 fn name(&self) -> &str {
1303 "mixed-workflow"
1304 }
1305
1306 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
1307 Box::pin(async move {
1308 ctx.shell("build", ShellConfig::new("echo built")).await?;
1309 let op = FakeGitlabOp {
1310 project_id: 456,
1311 title: "Deploy done".to_string(),
1312 };
1313 let result = ctx.operation("notify-gitlab", &op).await?;
1314 assert_eq!(result.output["issue_id"], 42);
1315 Ok(())
1316 })
1317 }
1318 }
1319
1320 #[tokio::test]
1321 async fn operation_step_happy_path() {
1322 let mut engine = create_test_engine();
1323 engine.register(OperationWorkflow).unwrap();
1324
1325 let run = engine
1326 .run_handler("operation-workflow", TriggerKind::Manual, json!({}))
1327 .await
1328 .unwrap();
1329
1330 assert_eq!(run.status.state, RunStatus::Completed);
1331
1332 let steps = engine.store().list_steps(run.id).await.unwrap();
1333
1334 assert_eq!(steps.len(), 1);
1335 assert_eq!(steps[0].name, "create-issue");
1336 assert_eq!(steps[0].kind, StepKind::Custom("gitlab".to_string()));
1337 assert_eq!(
1338 steps[0].status.state,
1339 ironflow_store::models::StepStatus::Completed
1340 );
1341
1342 let output = steps[0].output.as_ref().unwrap();
1343 assert_eq!(output["issue_id"], 42);
1344 assert_eq!(output["project_id"], 123);
1345
1346 let input = steps[0].input.as_ref().unwrap();
1347 assert_eq!(input["project_id"], 123);
1348 assert_eq!(input["title"], "Bug report");
1349 }
1350
1351 #[tokio::test]
1352 async fn operation_step_failure_marks_run_failed() {
1353 let mut engine = create_test_engine();
1354 engine.register(FailingOperationWorkflow).unwrap();
1355
1356 let result = engine
1357 .run_handler("failing-operation-workflow", TriggerKind::Manual, json!({}))
1358 .await;
1359
1360 assert!(result.is_err());
1361 }
1362
1363 #[tokio::test]
1364 async fn operation_mixed_with_shell_steps() {
1365 let mut engine = create_test_engine();
1366 engine.register(MixedWorkflow).unwrap();
1367
1368 let run = engine
1369 .run_handler("mixed-workflow", TriggerKind::Manual, json!({}))
1370 .await
1371 .unwrap();
1372
1373 assert_eq!(run.status.state, RunStatus::Completed);
1374
1375 let steps = engine.store().list_steps(run.id).await.unwrap();
1376
1377 assert_eq!(steps.len(), 2);
1378 assert_eq!(steps[0].kind, StepKind::Shell);
1379 assert_eq!(steps[1].kind, StepKind::Custom("gitlab".to_string()));
1380 assert_eq!(steps[0].position, 0);
1381 assert_eq!(steps[1].position, 1);
1382 }
1383
1384 use crate::config::ApprovalConfig;
1389
1390 struct SingleApprovalWorkflow;
1391
1392 impl WorkflowHandler for SingleApprovalWorkflow {
1393 fn name(&self) -> &str {
1394 "single-approval"
1395 }
1396
1397 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
1398 Box::pin(async move {
1399 ctx.shell("build", ShellConfig::new("echo built")).await?;
1400 ctx.approval("gate", ApprovalConfig::new("OK?")).await?;
1401 ctx.shell("deploy", ShellConfig::new("echo deployed"))
1402 .await?;
1403 Ok(())
1404 })
1405 }
1406 }
1407
1408 struct DoubleApprovalWorkflow;
1409
1410 impl WorkflowHandler for DoubleApprovalWorkflow {
1411 fn name(&self) -> &str {
1412 "double-approval"
1413 }
1414
1415 fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
1416 Box::pin(async move {
1417 ctx.shell("build", ShellConfig::new("echo built")).await?;
1418 ctx.approval("staging-gate", ApprovalConfig::new("Deploy staging?"))
1419 .await?;
1420 ctx.shell("deploy-staging", ShellConfig::new("echo staging"))
1421 .await?;
1422 ctx.approval("prod-gate", ApprovalConfig::new("Deploy prod?"))
1423 .await?;
1424 ctx.shell("deploy-prod", ShellConfig::new("echo prod"))
1425 .await?;
1426 Ok(())
1427 })
1428 }
1429 }
1430
1431 #[tokio::test]
1432 async fn approval_pauses_run() {
1433 let mut engine = create_test_engine();
1434 engine.register(SingleApprovalWorkflow).unwrap();
1435
1436 let run = engine
1437 .run_handler("single-approval", TriggerKind::Manual, json!({}))
1438 .await
1439 .unwrap();
1440
1441 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
1442
1443 let steps = engine.store().list_steps(run.id).await.unwrap();
1444 assert_eq!(steps.len(), 2); assert_eq!(steps[0].kind, StepKind::Shell);
1446 assert_eq!(steps[0].status.state, StepStatus::Completed);
1447 assert_eq!(steps[1].kind, StepKind::Approval);
1448 assert_eq!(steps[1].status.state, StepStatus::AwaitingApproval);
1449 }
1450
1451 #[tokio::test]
1452 async fn approval_resume_completes_run() {
1453 let mut engine = create_test_engine();
1454 engine.register(SingleApprovalWorkflow).unwrap();
1455
1456 let run = engine
1458 .run_handler("single-approval", TriggerKind::Manual, json!({}))
1459 .await
1460 .unwrap();
1461 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
1462
1463 engine
1465 .store()
1466 .update_run_status(run.id, RunStatus::Running)
1467 .await
1468 .unwrap();
1469
1470 let resumed = engine.resume_run(run.id).await.unwrap();
1472 assert_eq!(resumed.status.state, RunStatus::Completed);
1473
1474 let steps = engine.store().list_steps(run.id).await.unwrap();
1475 assert_eq!(steps.len(), 3); assert_eq!(steps[0].name, "build");
1477 assert_eq!(steps[0].status.state, StepStatus::Completed);
1478 assert_eq!(steps[1].name, "gate");
1479 assert_eq!(steps[1].kind, StepKind::Approval);
1480 assert_eq!(steps[1].status.state, StepStatus::Completed);
1481 assert_eq!(steps[2].name, "deploy");
1482 assert_eq!(steps[2].status.state, StepStatus::Completed);
1483 }
1484
1485 #[tokio::test]
1486 async fn double_approval_two_resumes() {
1487 let mut engine = create_test_engine();
1488 engine.register(DoubleApprovalWorkflow).unwrap();
1489
1490 let run = engine
1492 .run_handler("double-approval", TriggerKind::Manual, json!({}))
1493 .await
1494 .unwrap();
1495 assert_eq!(run.status.state, RunStatus::AwaitingApproval);
1496
1497 let steps = engine.store().list_steps(run.id).await.unwrap();
1498 assert_eq!(steps.len(), 2); engine
1502 .store()
1503 .update_run_status(run.id, RunStatus::Running)
1504 .await
1505 .unwrap();
1506
1507 let resumed = engine.resume_run(run.id).await.unwrap();
1508 assert_eq!(resumed.status.state, RunStatus::AwaitingApproval);
1509
1510 let steps = engine.store().list_steps(run.id).await.unwrap();
1511 assert_eq!(steps.len(), 4); engine
1515 .store()
1516 .update_run_status(run.id, RunStatus::Running)
1517 .await
1518 .unwrap();
1519
1520 let final_run = engine.resume_run(run.id).await.unwrap();
1521 assert_eq!(final_run.status.state, RunStatus::Completed);
1522
1523 let steps = engine.store().list_steps(run.id).await.unwrap();
1524 assert_eq!(steps.len(), 5);
1525 assert_eq!(steps[0].name, "build");
1526 assert_eq!(steps[1].name, "staging-gate");
1527 assert_eq!(steps[2].name, "deploy-staging");
1528 assert_eq!(steps[3].name, "prod-gate");
1529 assert_eq!(steps[4].name, "deploy-prod");
1530
1531 for step in &steps {
1532 assert_eq!(step.status.state, StepStatus::Completed);
1533 }
1534 }
1535
1536 use ironflow_store::models::{NewStep, StepUpdate};
1541
1542 async fn create_step_with_status(
1543 store: &Arc<dyn Store>,
1544 run_id: Uuid,
1545 name: &str,
1546 position: u32,
1547 status: StepStatus,
1548 ) -> ironflow_store::models::Step {
1549 let step = store
1550 .create_step(NewStep {
1551 run_id,
1552 name: name.to_string(),
1553 kind: StepKind::Shell,
1554 position,
1555 input: None,
1556 })
1557 .await
1558 .unwrap();
1559
1560 match status {
1561 StepStatus::Pending => {}
1562 StepStatus::Running => {
1563 store
1564 .update_step(
1565 step.id,
1566 StepUpdate {
1567 status: Some(StepStatus::Running),
1568 ..StepUpdate::default()
1569 },
1570 )
1571 .await
1572 .unwrap();
1573 }
1574 StepStatus::Completed => {
1575 store
1576 .update_step(
1577 step.id,
1578 StepUpdate {
1579 status: Some(StepStatus::Running),
1580 ..StepUpdate::default()
1581 },
1582 )
1583 .await
1584 .unwrap();
1585 store
1586 .update_step(
1587 step.id,
1588 StepUpdate {
1589 status: Some(StepStatus::Completed),
1590 ..StepUpdate::default()
1591 },
1592 )
1593 .await
1594 .unwrap();
1595 }
1596 StepStatus::AwaitingApproval => {
1597 store
1598 .update_step(
1599 step.id,
1600 StepUpdate {
1601 status: Some(StepStatus::Running),
1602 ..StepUpdate::default()
1603 },
1604 )
1605 .await
1606 .unwrap();
1607 store
1608 .update_step(
1609 step.id,
1610 StepUpdate {
1611 status: Some(StepStatus::AwaitingApproval),
1612 ..StepUpdate::default()
1613 },
1614 )
1615 .await
1616 .unwrap();
1617 }
1618 _ => panic!("unsupported status for test helper: {status}"),
1619 }
1620
1621 store.get_step(step.id).await.unwrap().unwrap()
1622 }
1623
1624 #[tokio::test]
1625 async fn fail_orphaned_steps_marks_running_as_failed() {
1626 let engine = create_test_engine();
1627 let run = engine
1628 .store()
1629 .create_run(NewRun {
1630 workflow_name: "test".to_string(),
1631 trigger: TriggerKind::Manual,
1632 payload: json!({}),
1633 max_retries: 0,
1634 handler_version: None,
1635 labels: HashMap::new(),
1636 scheduled_at: None,
1637 })
1638 .await
1639 .unwrap();
1640
1641 let step = create_step_with_status(
1642 engine.store(),
1643 run.id,
1644 "running-step",
1645 0,
1646 StepStatus::Running,
1647 )
1648 .await;
1649
1650 engine
1651 .fail_orphaned_steps(run.id, "parent run timed out")
1652 .await
1653 .unwrap();
1654
1655 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
1656 assert_eq!(updated.status.state, StepStatus::Failed);
1657 assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
1658 assert!(updated.completed_at.is_some());
1659 }
1660
1661 #[tokio::test]
1662 async fn fail_orphaned_steps_marks_pending_as_skipped() {
1663 let engine = create_test_engine();
1664 let run = engine
1665 .store()
1666 .create_run(NewRun {
1667 workflow_name: "test".to_string(),
1668 trigger: TriggerKind::Manual,
1669 payload: json!({}),
1670 max_retries: 0,
1671 handler_version: None,
1672 labels: HashMap::new(),
1673 scheduled_at: None,
1674 })
1675 .await
1676 .unwrap();
1677
1678 let step = create_step_with_status(
1679 engine.store(),
1680 run.id,
1681 "pending-step",
1682 0,
1683 StepStatus::Pending,
1684 )
1685 .await;
1686
1687 engine
1688 .fail_orphaned_steps(run.id, "parent run timed out")
1689 .await
1690 .unwrap();
1691
1692 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
1693 assert_eq!(updated.status.state, StepStatus::Skipped);
1694 assert!(updated.error.is_none());
1695 assert!(updated.completed_at.is_some());
1696 }
1697
1698 #[tokio::test]
1699 async fn fail_orphaned_steps_marks_awaiting_approval_as_failed() {
1700 let engine = create_test_engine();
1701 let run = engine
1702 .store()
1703 .create_run(NewRun {
1704 workflow_name: "test".to_string(),
1705 trigger: TriggerKind::Manual,
1706 payload: json!({}),
1707 max_retries: 0,
1708 handler_version: None,
1709 labels: HashMap::new(),
1710 scheduled_at: None,
1711 })
1712 .await
1713 .unwrap();
1714
1715 let step = create_step_with_status(
1716 engine.store(),
1717 run.id,
1718 "approval-step",
1719 0,
1720 StepStatus::AwaitingApproval,
1721 )
1722 .await;
1723
1724 engine
1725 .fail_orphaned_steps(run.id, "parent run timed out")
1726 .await
1727 .unwrap();
1728
1729 let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
1730 assert_eq!(updated.status.state, StepStatus::Failed);
1731 assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
1732 assert!(updated.completed_at.is_some());
1733 }
1734
1735 #[tokio::test]
1736 async fn fail_orphaned_steps_skips_terminal_steps() {
1737 let engine = create_test_engine();
1738 let run = engine
1739 .store()
1740 .create_run(NewRun {
1741 workflow_name: "test".to_string(),
1742 trigger: TriggerKind::Manual,
1743 payload: json!({}),
1744 max_retries: 0,
1745 handler_version: None,
1746 labels: HashMap::new(),
1747 scheduled_at: None,
1748 })
1749 .await
1750 .unwrap();
1751
1752 let completed_step =
1753 create_step_with_status(engine.store(), run.id, "done", 0, StepStatus::Completed).await;
1754 let running_step =
1755 create_step_with_status(engine.store(), run.id, "in-flight", 1, StepStatus::Running)
1756 .await;
1757
1758 engine
1759 .fail_orphaned_steps(run.id, "parent run timed out")
1760 .await
1761 .unwrap();
1762
1763 let completed = engine
1764 .store()
1765 .get_step(completed_step.id)
1766 .await
1767 .unwrap()
1768 .unwrap();
1769 assert_eq!(completed.status.state, StepStatus::Completed);
1770
1771 let failed = engine
1772 .store()
1773 .get_step(running_step.id)
1774 .await
1775 .unwrap()
1776 .unwrap();
1777 assert_eq!(failed.status.state, StepStatus::Failed);
1778 }
1779
1780 #[tokio::test]
1781 async fn fail_orphaned_steps_mixed_states() {
1782 let engine = create_test_engine();
1783 let run = engine
1784 .store()
1785 .create_run(NewRun {
1786 workflow_name: "test".to_string(),
1787 trigger: TriggerKind::Manual,
1788 payload: json!({}),
1789 max_retries: 0,
1790 handler_version: None,
1791 labels: HashMap::new(),
1792 scheduled_at: None,
1793 })
1794 .await
1795 .unwrap();
1796
1797 let s_completed =
1798 create_step_with_status(engine.store(), run.id, "step-1", 0, StepStatus::Completed)
1799 .await;
1800 let s_running =
1801 create_step_with_status(engine.store(), run.id, "step-2", 1, StepStatus::Running).await;
1802 let s_pending =
1803 create_step_with_status(engine.store(), run.id, "step-3", 2, StepStatus::Pending).await;
1804
1805 engine.fail_orphaned_steps(run.id, "timeout").await.unwrap();
1806
1807 let r_completed = engine
1808 .store()
1809 .get_step(s_completed.id)
1810 .await
1811 .unwrap()
1812 .unwrap();
1813 assert_eq!(r_completed.status.state, StepStatus::Completed);
1814
1815 let r_running = engine
1816 .store()
1817 .get_step(s_running.id)
1818 .await
1819 .unwrap()
1820 .unwrap();
1821 assert_eq!(r_running.status.state, StepStatus::Failed);
1822 assert_eq!(r_running.error.as_deref(), Some("timeout"));
1823
1824 let r_pending = engine
1825 .store()
1826 .get_step(s_pending.id)
1827 .await
1828 .unwrap()
1829 .unwrap();
1830 assert_eq!(r_pending.status.state, StepStatus::Skipped);
1831 assert!(r_pending.error.is_none());
1832 }
1833
1834 #[tokio::test]
1835 async fn fail_orphaned_steps_no_steps_is_noop() {
1836 let engine = create_test_engine();
1837 let run = engine
1838 .store()
1839 .create_run(NewRun {
1840 workflow_name: "test".to_string(),
1841 trigger: TriggerKind::Manual,
1842 payload: json!({}),
1843 max_retries: 0,
1844 handler_version: None,
1845 labels: HashMap::new(),
1846 scheduled_at: None,
1847 })
1848 .await
1849 .unwrap();
1850
1851 let result = engine.fail_orphaned_steps(run.id, "timeout").await;
1852 assert!(result.is_ok());
1853 }
1854
1855 #[tokio::test]
1856 async fn fail_orphaned_steps_preserves_existing_error() {
1857 let engine = create_test_engine();
1858 let run = engine
1859 .store()
1860 .create_run(NewRun {
1861 workflow_name: "test".to_string(),
1862 trigger: TriggerKind::Manual,
1863 payload: json!({}),
1864 max_retries: 0,
1865 handler_version: None,
1866 labels: HashMap::new(),
1867 scheduled_at: None,
1868 })
1869 .await
1870 .unwrap();
1871
1872 let step_with_error = create_step_with_status(
1873 engine.store(),
1874 run.id,
1875 "already-errored",
1876 0,
1877 StepStatus::Running,
1878 )
1879 .await;
1880
1881 engine
1882 .store()
1883 .update_step(
1884 step_with_error.id,
1885 StepUpdate {
1886 error: Some("real error from provider".to_string()),
1887 ..StepUpdate::default()
1888 },
1889 )
1890 .await
1891 .unwrap();
1892
1893 let step_no_error = create_step_with_status(
1894 engine.store(),
1895 run.id,
1896 "no-error-yet",
1897 1,
1898 StepStatus::Running,
1899 )
1900 .await;
1901
1902 engine
1903 .fail_orphaned_steps(run.id, "parent run failed")
1904 .await
1905 .unwrap();
1906
1907 let updated_with = engine
1908 .store()
1909 .get_step(step_with_error.id)
1910 .await
1911 .unwrap()
1912 .unwrap();
1913 assert_eq!(updated_with.status.state, StepStatus::Failed);
1914 assert_eq!(
1915 updated_with.error.as_deref(),
1916 Some("real error from provider"),
1917 );
1918
1919 let updated_without = engine
1920 .store()
1921 .get_step(step_no_error.id)
1922 .await
1923 .unwrap()
1924 .unwrap();
1925 assert_eq!(updated_without.status.state, StepStatus::Failed);
1926 assert_eq!(updated_without.error.as_deref(), Some("parent run failed"),);
1927 }
1928}