1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
//! Explicitly skipped step for [`WorkflowContext`].
use chrono::Utc;
use serde_json::json;
use tracing::info;
use ironflow_store::models::{
NewStep, NewStepDependency, StepKind, StepStatus, StepUpdate, step_trace_id,
};
use crate::context::WorkflowContext;
use crate::error::EngineError;
use crate::plan::{ConditionResult, lock_plan};
impl WorkflowContext {
/// Record a step as explicitly skipped.
///
/// Use this inside an `if`/`else` branch when a step should not execute
/// but must still appear in the DAG and timeline with its reason.
///
/// The step is created directly in [`StepStatus::Skipped`] state and the
/// reason is stored in the output as `{"reason": "..."}`.
///
/// On resume, a skip already recorded in a prior execution of the current
/// attempt is replayed instead of creating a second `Skipped` step.
///
/// # Errors
///
/// Returns [`EngineError`] if the store fails.
///
/// # Examples
///
/// ```no_run
/// use ironflow_engine::context::WorkflowContext;
/// use ironflow_engine::error::EngineError;
///
/// # async fn example(ctx: &mut WorkflowContext) -> Result<(), EngineError> {
/// let tests_passed = false;
/// if tests_passed {
/// // ctx.shell("deploy", ...).await?;
/// } else {
/// ctx.skip("deploy", "tests failed").await?;
/// }
/// # Ok(())
/// # }
/// ```
pub async fn skip(&mut self, name: &str, reason: &str) -> Result<(), EngineError> {
// Plan mode: the skip and its reason become the step's condition.
if let Some(plan) = self.plan().cloned() {
self.position += 1;
let mut recorder = lock_plan(&plan);
recorder.set_condition(ConditionResult::Skipped {
reason: reason.to_string(),
});
if recorder.record(
name,
StepKind::Custom("skip".to_string()),
&self.workflow_name,
None,
) {
recorder.set_last(vec![name.to_string()]);
}
return Ok(());
}
let position = self.position;
self.position += 1;
if let Some(existing) = self.replay_steps.get(&position)
&& existing.name == name
&& existing.kind == StepKind::Custom("skip".to_string())
&& existing.status.state == StepStatus::Skipped
{
self.last_step_ids = vec![existing.id];
info!(
run_id = %self.run_id,
step = %name,
position,
"step replayed from previous execution"
);
return Ok(());
}
let trace_id = step_trace_id(self.run_id, name, position);
let step = self
.store
.create_step(NewStep {
run_id: self.run_id,
trace_id,
name: name.to_string(),
kind: StepKind::Custom("skip".to_string()),
position,
input: None,
is_error_handler: false,
})
.await?;
if !self.last_step_ids.is_empty() {
let deps: Vec<NewStepDependency> = self
.last_step_ids
.iter()
.map(|&depends_on| NewStepDependency {
step_id: step.id,
depends_on,
})
.collect();
self.store.create_step_dependencies(deps).await?;
}
let now = Utc::now();
self.store
.update_step(
step.id,
StepUpdate {
status: Some(StepStatus::Skipped),
output: Some(json!({"reason": reason})),
completed_at: Some(now),
..StepUpdate::default()
},
)
.await?;
self.last_step_ids = vec![step.id];
info!(
run_id = %self.run_id,
step = %name,
reason,
"step skipped"
);
Ok(())
}
}