ironflow_engine/context/steps/operation.rs
1//! Custom operation step for [`WorkflowContext`].
2
3use std::time::Instant;
4
5use chrono::Utc;
6use rust_decimal::Decimal;
7use tracing::{error, info};
8
9use ironflow_store::models::{NewStep, StepKind, StepStatus, StepUpdate, step_trace_id};
10
11use crate::context::WorkflowContext;
12use crate::error::EngineError;
13use crate::executor::{StepArtifacts, StepOutput};
14use crate::operation::Operation;
15use crate::plan::{lock_plan, planned_custom_output};
16
17impl WorkflowContext {
18 /// Execute a custom operation step.
19 ///
20 /// Runs a user-defined [`Operation`] with full step lifecycle management:
21 /// creates the step record, transitions to Running, executes the operation,
22 /// persists the output and duration, and marks the step Completed or Failed.
23 ///
24 /// The operation's [`kind()`](Operation::kind) is stored as
25 /// [`StepKind::Custom`].
26 ///
27 /// On resume, a step that already completed in a prior execution of the
28 /// current attempt is replayed from the store instead of calling
29 /// [`Operation::execute`] again.
30 ///
31 /// # Errors
32 ///
33 /// Returns [`EngineError`] if the operation fails or the store errors.
34 /// Returns [`EngineError::ReplayDivergence`] when the step recorded at
35 /// this position has a different name or kind.
36 ///
37 /// # Examples
38 ///
39 /// ```no_run
40 /// use async_trait::async_trait;
41 /// use ironflow_engine::context::WorkflowContext;
42 /// use ironflow_engine::operation::{Operation, OperationContext};
43 /// use ironflow_core::error::OperationError;
44 /// use ironflow_engine::error::EngineError;
45 /// use serde_json::{Value, json};
46 ///
47 /// struct MyOp;
48 /// #[async_trait]
49 /// impl Operation for MyOp {
50 /// fn kind(&self) -> &str { "my-service" }
51 /// async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
52 /// Ok(json!({"ok": true}))
53 /// }
54 /// }
55 ///
56 /// # async fn example(ctx: &mut WorkflowContext) -> Result<(), EngineError> {
57 /// let result = ctx.operation("call-service", &MyOp).await?;
58 /// println!("output: {}", result.output);
59 /// # Ok(())
60 /// # }
61 /// ```
62 pub async fn operation(
63 &mut self,
64 name: &str,
65 op: &dyn Operation,
66 ) -> Result<StepOutput, EngineError> {
67 let kind = StepKind::Custom(op.kind().to_string());
68
69 // Plan mode: `op.execute` is never called, so no third-party API is
70 // touched while planning.
71 if let Some(plan) = self.plan().cloned() {
72 self.position += 1;
73 let mut recorder = lock_plan(&plan);
74 let estimate = recorder.estimate_for(name);
75 if recorder.record(name, kind, &self.workflow_name, None) {
76 recorder.set_last(vec![name.to_string()]);
77 }
78 let mut output = planned_custom_output(estimate);
79 output.artifacts = StepArtifacts::new(name, None, &[]);
80 return Ok(output);
81 }
82
83 let position = self.position;
84 self.position += 1;
85
86 // Replay: if this step already completed in a prior execution of the
87 // current attempt, return its cached output without calling
88 // `op.execute` or creating a new step.
89 if let Some(mut output) = self.try_replay_step(position, name, &kind)? {
90 let step_id = self.last_step_ids.last().copied();
91 output.artifacts = StepArtifacts::new(name, step_id, &[]);
92 return Ok(output);
93 }
94
95 let trace_id = step_trace_id(self.run_id, name, position);
96 let step = self
97 .store
98 .create_step(NewStep {
99 run_id: self.run_id,
100 trace_id,
101 name: name.to_string(),
102 kind,
103 position,
104 input: op.input(),
105 is_error_handler: false,
106 })
107 .await?;
108
109 self.start_step(step.id, Utc::now()).await?;
110
111 let start = Instant::now();
112
113 let op_ctx = self.ensure_operation_ctx();
114
115 match op.execute(op_ctx).await {
116 Ok(output_value) => {
117 let duration_ms = start.elapsed().as_millis() as u64;
118 self.total_duration_ms += duration_ms;
119
120 let completed_at = Utc::now();
121 self.store
122 .update_step(
123 step.id,
124 StepUpdate {
125 status: Some(StepStatus::Completed),
126 output: Some(output_value.clone()),
127 duration_ms: Some(duration_ms),
128 cost_usd: Some(Decimal::ZERO),
129 completed_at: Some(completed_at),
130 ..StepUpdate::default()
131 },
132 )
133 .await?;
134
135 info!(
136 run_id = %self.run_id,
137 step = %name,
138 kind = op.kind(),
139 duration_ms,
140 "operation step completed"
141 );
142
143 self.last_step_ids = vec![step.id];
144
145 Ok(StepOutput {
146 output: output_value,
147 duration_ms,
148 cost_usd: Decimal::ZERO,
149 input_tokens: None,
150 cache_read_input_tokens: None,
151 cache_creation_input_tokens: None,
152 output_tokens: None,
153 model: None,
154 debug_messages: None,
155 // An operation declares no output; its record is what
156 // `put_artifact` attaches bytes to.
157 artifacts: StepArtifacts::new(name, Some(step.id), &[]),
158 })
159 }
160 Err(err) => {
161 let completed_at = Utc::now();
162 let engine_err = EngineError::Operation(err);
163 if let Err(store_err) = self
164 .store
165 .update_step(
166 step.id,
167 StepUpdate {
168 status: Some(StepStatus::Failed),
169 error: Some(engine_err.to_string()),
170 completed_at: Some(completed_at),
171 ..StepUpdate::default()
172 },
173 )
174 .await
175 {
176 error!(step_id = %step.id, error = %store_err, "failed to persist step failure");
177 }
178
179 Err(engine_err)
180 }
181 }
182 }
183}