ironflow_engine/context/steps/
skip.rs1use chrono::Utc;
4use serde_json::json;
5use tracing::info;
6
7use ironflow_store::models::{
8 NewStep, NewStepDependency, StepKind, StepStatus, StepUpdate, step_trace_id,
9};
10
11use crate::context::WorkflowContext;
12use crate::error::EngineError;
13use crate::plan::{ConditionResult, lock_plan};
14
15impl WorkflowContext {
16 pub async fn skip(&mut self, name: &str, reason: &str) -> Result<(), EngineError> {
48 if let Some(plan) = self.plan().cloned() {
50 self.position += 1;
51 let mut recorder = lock_plan(&plan);
52 recorder.set_condition(ConditionResult::Skipped {
53 reason: reason.to_string(),
54 });
55 if recorder.record(
56 name,
57 StepKind::Custom("skip".to_string()),
58 &self.workflow_name,
59 None,
60 ) {
61 recorder.set_last(vec![name.to_string()]);
62 }
63 return Ok(());
64 }
65
66 let position = self.position;
67 self.position += 1;
68
69 if let Some(existing) = self.replay_steps.get(&position)
70 && existing.name == name
71 && existing.kind == StepKind::Custom("skip".to_string())
72 && existing.status.state == StepStatus::Skipped
73 {
74 self.last_step_ids = vec![existing.id];
75 info!(
76 run_id = %self.run_id,
77 step = %name,
78 position,
79 "step replayed from previous execution"
80 );
81 return Ok(());
82 }
83
84 let trace_id = step_trace_id(self.run_id, name, position);
85 let step = self
86 .store
87 .create_step(NewStep {
88 run_id: self.run_id,
89 trace_id,
90 name: name.to_string(),
91 kind: StepKind::Custom("skip".to_string()),
92 position,
93 input: None,
94 is_error_handler: false,
95 })
96 .await?;
97
98 if !self.last_step_ids.is_empty() {
99 let deps: Vec<NewStepDependency> = self
100 .last_step_ids
101 .iter()
102 .map(|&depends_on| NewStepDependency {
103 step_id: step.id,
104 depends_on,
105 })
106 .collect();
107 self.store.create_step_dependencies(deps).await?;
108 }
109
110 let now = Utc::now();
111 self.store
112 .update_step(
113 step.id,
114 StepUpdate {
115 status: Some(StepStatus::Skipped),
116 output: Some(json!({"reason": reason})),
117 completed_at: Some(now),
118 ..StepUpdate::default()
119 },
120 )
121 .await?;
122
123 self.last_step_ids = vec![step.id];
124
125 info!(
126 run_id = %self.run_id,
127 step = %name,
128 reason,
129 "step skipped"
130 );
131
132 Ok(())
133 }
134}