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