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 /// # 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 let mut output = planned_custom_output(estimate);
73 output.artifacts = StepArtifacts::new(name, None, &[]);
74 return Ok(output);
75 }
76
77 let position = self.position;
78 self.position += 1;
79
80 let trace_id = step_trace_id(self.run_id, name, position);
81 let step = self
82 .store
83 .create_step(NewStep {
84 run_id: self.run_id,
85 trace_id,
86 name: name.to_string(),
87 kind,
88 position,
89 input: op.input(),
90 is_error_handler: false,
91 })
92 .await?;
93
94 self.start_step(step.id, Utc::now()).await?;
95
96 let start = Instant::now();
97
98 let op_ctx = self.ensure_operation_ctx();
99
100 match op.execute(op_ctx).await {
101 Ok(output_value) => {
102 let duration_ms = start.elapsed().as_millis() as u64;
103 self.total_duration_ms += duration_ms;
104
105 let completed_at = Utc::now();
106 self.store
107 .update_step(
108 step.id,
109 StepUpdate {
110 status: Some(StepStatus::Completed),
111 output: Some(output_value.clone()),
112 duration_ms: Some(duration_ms),
113 cost_usd: Some(Decimal::ZERO),
114 completed_at: Some(completed_at),
115 ..StepUpdate::default()
116 },
117 )
118 .await?;
119
120 info!(
121 run_id = %self.run_id,
122 step = %name,
123 kind = op.kind(),
124 duration_ms,
125 "operation step completed"
126 );
127
128 self.last_step_ids = vec![step.id];
129
130 Ok(StepOutput {
131 output: output_value,
132 duration_ms,
133 cost_usd: Decimal::ZERO,
134 input_tokens: None,
135 cache_read_input_tokens: None,
136 cache_creation_input_tokens: None,
137 output_tokens: None,
138 model: None,
139 debug_messages: None,
140 // An operation declares no output; its record is what
141 // `put_artifact` attaches bytes to.
142 artifacts: StepArtifacts::new(name, Some(step.id), &[]),
143 })
144 }
145 Err(err) => {
146 let completed_at = Utc::now();
147 let engine_err = EngineError::Operation(err);
148 if let Err(store_err) = self
149 .store
150 .update_step(
151 step.id,
152 StepUpdate {
153 status: Some(StepStatus::Failed),
154 error: Some(engine_err.to_string()),
155 completed_at: Some(completed_at),
156 ..StepUpdate::default()
157 },
158 )
159 .await
160 {
161 error!(step_id = %step.id, error = %store_err, "failed to persist step failure");
162 }
163
164 Err(engine_err)
165 }
166 }
167 }
168}