Skip to main content

ironflow_engine/context/steps/
skip.rs

1//! Explicitly skipped step for [`WorkflowContext`].
2
3use 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    /// Record a step as explicitly skipped.
17    ///
18    /// Use this inside an `if`/`else` branch when a step should not execute
19    /// but must still appear in the DAG and timeline with its reason.
20    ///
21    /// The step is created directly in [`StepStatus::Skipped`] state and the
22    /// reason is stored in the output as `{"reason": "..."}`.
23    ///
24    /// On resume, a skip already recorded in a prior execution of the current
25    /// attempt is replayed instead of creating a second `Skipped` step.
26    ///
27    /// # Errors
28    ///
29    /// Returns [`EngineError`] if the store fails.
30    ///
31    /// # Examples
32    ///
33    /// ```no_run
34    /// use ironflow_engine::context::WorkflowContext;
35    /// use ironflow_engine::error::EngineError;
36    ///
37    /// # async fn example(ctx: &mut WorkflowContext) -> Result<(), EngineError> {
38    /// let tests_passed = false;
39    /// if tests_passed {
40    ///     // ctx.shell("deploy", ...).await?;
41    /// } else {
42    ///     ctx.skip("deploy", "tests failed").await?;
43    /// }
44    /// # Ok(())
45    /// # }
46    /// ```
47    pub async fn skip(&mut self, name: &str, reason: &str) -> Result<(), EngineError> {
48        // Plan mode: the skip and its reason become the step's condition.
49        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}