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 ///
35 /// # Examples
36 ///
37 /// ```no_run
38 /// use async_trait::async_trait;
39 /// use ironflow_engine::context::WorkflowContext;
40 /// use ironflow_engine::operation::{Operation, OperationContext};
41 /// use ironflow_core::error::OperationError;
42 /// use ironflow_engine::error::EngineError;
43 /// use serde_json::{Value, json};
44 ///
45 /// struct MyOp;
46 /// #[async_trait]
47 /// impl Operation for MyOp {
48 /// fn kind(&self) -> &str { "my-service" }
49 /// async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
50 /// Ok(json!({"ok": true}))
51 /// }
52 /// }
53 ///
54 /// # async fn example(ctx: &mut WorkflowContext) -> Result<(), EngineError> {
55 /// let result = ctx.operation("call-service", &MyOp).await?;
56 /// println!("output: {}", result.output);
57 /// # Ok(())
58 /// # }
59 /// ```
60 pub async fn operation(
61 &mut self,
62 name: &str,
63 op: &dyn Operation,
64 ) -> Result<StepOutput, EngineError> {
65 let kind = StepKind::Custom(op.kind().to_string());
66
67 // Plan mode: `op.execute` is never called, so no third-party API is
68 // touched while planning.
69 if let Some(plan) = self.plan().cloned() {
70 self.position += 1;
71 let mut recorder = lock_plan(&plan);
72 let estimate = recorder.estimate_for(name);
73 if recorder.record(name, kind, &self.workflow_name, None) {
74 recorder.set_last(vec![name.to_string()]);
75 }
76 let mut output = planned_custom_output(estimate);
77 output.artifacts = StepArtifacts::new(name, None, &[]);
78 return Ok(output);
79 }
80
81 let position = self.position;
82 self.position += 1;
83
84 // Replay: if this step already completed in a prior execution of the
85 // current attempt, return its cached output without calling
86 // `op.execute` or creating a new step.
87 if let Some(mut output) = self.try_replay_step(position) {
88 let step_id = self.last_step_ids.last().copied();
89 output.artifacts = StepArtifacts::new(name, step_id, &[]);
90 return Ok(output);
91 }
92
93 let trace_id = step_trace_id(self.run_id, name, position);
94 let step = self
95 .store
96 .create_step(NewStep {
97 run_id: self.run_id,
98 trace_id,
99 name: name.to_string(),
100 kind,
101 position,
102 input: op.input(),
103 is_error_handler: false,
104 })
105 .await?;
106
107 self.start_step(step.id, Utc::now()).await?;
108
109 let start = Instant::now();
110
111 let op_ctx = self.ensure_operation_ctx();
112
113 match op.execute(op_ctx).await {
114 Ok(output_value) => {
115 let duration_ms = start.elapsed().as_millis() as u64;
116 self.total_duration_ms += duration_ms;
117
118 let completed_at = Utc::now();
119 self.store
120 .update_step(
121 step.id,
122 StepUpdate {
123 status: Some(StepStatus::Completed),
124 output: Some(output_value.clone()),
125 duration_ms: Some(duration_ms),
126 cost_usd: Some(Decimal::ZERO),
127 completed_at: Some(completed_at),
128 ..StepUpdate::default()
129 },
130 )
131 .await?;
132
133 info!(
134 run_id = %self.run_id,
135 step = %name,
136 kind = op.kind(),
137 duration_ms,
138 "operation step completed"
139 );
140
141 self.last_step_ids = vec![step.id];
142
143 Ok(StepOutput {
144 output: output_value,
145 duration_ms,
146 cost_usd: Decimal::ZERO,
147 input_tokens: None,
148 cache_read_input_tokens: None,
149 cache_creation_input_tokens: None,
150 output_tokens: None,
151 model: None,
152 debug_messages: None,
153 // An operation declares no output; its record is what
154 // `put_artifact` attaches bytes to.
155 artifacts: StepArtifacts::new(name, Some(step.id), &[]),
156 })
157 }
158 Err(err) => {
159 let completed_at = Utc::now();
160 let engine_err = EngineError::Operation(err);
161 if let Err(store_err) = self
162 .store
163 .update_step(
164 step.id,
165 StepUpdate {
166 status: Some(StepStatus::Failed),
167 error: Some(engine_err.to_string()),
168 completed_at: Some(completed_at),
169 ..StepUpdate::default()
170 },
171 )
172 .await
173 {
174 error!(step_id = %step.id, error = %store_err, "failed to persist step failure");
175 }
176
177 Err(engine_err)
178 }
179 }
180 }
181}