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;
13
14impl WorkflowContext {
15    /// Record a step as explicitly skipped.
16    ///
17    /// Use this inside an `if`/`else` branch when a step should not execute
18    /// but must still appear in the DAG and timeline with its reason.
19    ///
20    /// The step is created directly in [`StepStatus::Skipped`] state and the
21    /// reason is stored in the output as `{"reason": "..."}`.
22    ///
23    /// # Errors
24    ///
25    /// Returns [`EngineError`] if the store fails.
26    ///
27    /// # Examples
28    ///
29    /// ```no_run
30    /// use ironflow_engine::context::WorkflowContext;
31    /// use ironflow_engine::error::EngineError;
32    ///
33    /// # async fn example(ctx: &mut WorkflowContext) -> Result<(), EngineError> {
34    /// let tests_passed = false;
35    /// if tests_passed {
36    ///     // ctx.shell("deploy", ...).await?;
37    /// } else {
38    ///     ctx.skip("deploy", "tests failed").await?;
39    /// }
40    /// # Ok(())
41    /// # }
42    /// ```
43    pub async fn skip(&mut self, name: &str, reason: &str) -> Result<(), EngineError> {
44        let position = self.position;
45        self.position += 1;
46
47        let trace_id = step_trace_id(self.run_id, name, position);
48        let step = self
49            .store
50            .create_step(NewStep {
51                run_id: self.run_id,
52                trace_id,
53                name: name.to_string(),
54                kind: StepKind::Custom("skip".to_string()),
55                position,
56                input: None,
57                is_error_handler: false,
58            })
59            .await?;
60
61        if !self.last_step_ids.is_empty() {
62            let deps: Vec<NewStepDependency> = self
63                .last_step_ids
64                .iter()
65                .map(|&depends_on| NewStepDependency {
66                    step_id: step.id,
67                    depends_on,
68                })
69                .collect();
70            self.store.create_step_dependencies(deps).await?;
71        }
72
73        let now = Utc::now();
74        self.store
75            .update_step(
76                step.id,
77                StepUpdate {
78                    status: Some(StepStatus::Skipped),
79                    output: Some(json!({"reason": reason})),
80                    completed_at: Some(now),
81                    ..StepUpdate::default()
82                },
83            )
84            .await?;
85
86        self.last_step_ids = vec![step.id];
87
88        info!(
89            run_id = %self.run_id,
90            step = %name,
91            reason,
92            "step skipped"
93        );
94
95        Ok(())
96    }
97}