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;
15
16impl WorkflowContext {
17 /// Execute a custom operation step.
18 ///
19 /// Runs a user-defined [`Operation`] with full step lifecycle management:
20 /// creates the step record, transitions to Running, executes the operation,
21 /// persists the output and duration, and marks the step Completed or Failed.
22 ///
23 /// The operation's [`kind()`](Operation::kind) is stored as
24 /// [`StepKind::Custom`].
25 ///
26 /// # Errors
27 ///
28 /// Returns [`EngineError`] if the operation fails or the store errors.
29 ///
30 /// # Examples
31 ///
32 /// ```no_run
33 /// use async_trait::async_trait;
34 /// use ironflow_engine::context::WorkflowContext;
35 /// use ironflow_engine::operation::{Operation, OperationContext};
36 /// use ironflow_core::error::OperationError;
37 /// use ironflow_engine::error::EngineError;
38 /// use serde_json::{Value, json};
39 ///
40 /// struct MyOp;
41 /// #[async_trait]
42 /// impl Operation for MyOp {
43 /// fn kind(&self) -> &str { "my-service" }
44 /// async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
45 /// Ok(json!({"ok": true}))
46 /// }
47 /// }
48 ///
49 /// # async fn example(ctx: &mut WorkflowContext) -> Result<(), EngineError> {
50 /// let result = ctx.operation("call-service", &MyOp).await?;
51 /// println!("output: {}", result.output);
52 /// # Ok(())
53 /// # }
54 /// ```
55 pub async fn operation(
56 &mut self,
57 name: &str,
58 op: &dyn Operation,
59 ) -> Result<StepOutput, EngineError> {
60 let kind = StepKind::Custom(op.kind().to_string());
61 let position = self.position;
62 self.position += 1;
63
64 let trace_id = step_trace_id(self.run_id, name, position);
65 let step = self
66 .store
67 .create_step(NewStep {
68 run_id: self.run_id,
69 trace_id,
70 name: name.to_string(),
71 kind,
72 position,
73 input: op.input(),
74 is_error_handler: false,
75 })
76 .await?;
77
78 self.start_step(step.id, Utc::now()).await?;
79
80 let start = Instant::now();
81
82 let op_ctx = self.ensure_operation_ctx();
83
84 match op.execute(op_ctx).await {
85 Ok(output_value) => {
86 let duration_ms = start.elapsed().as_millis() as u64;
87 self.total_duration_ms += duration_ms;
88
89 let completed_at = Utc::now();
90 self.store
91 .update_step(
92 step.id,
93 StepUpdate {
94 status: Some(StepStatus::Completed),
95 output: Some(output_value.clone()),
96 duration_ms: Some(duration_ms),
97 cost_usd: Some(Decimal::ZERO),
98 completed_at: Some(completed_at),
99 ..StepUpdate::default()
100 },
101 )
102 .await?;
103
104 info!(
105 run_id = %self.run_id,
106 step = %name,
107 kind = op.kind(),
108 duration_ms,
109 "operation step completed"
110 );
111
112 self.last_step_ids = vec![step.id];
113
114 Ok(StepOutput {
115 output: output_value,
116 duration_ms,
117 cost_usd: Decimal::ZERO,
118 input_tokens: None,
119 output_tokens: None,
120 model: None,
121 debug_messages: None,
122 })
123 }
124 Err(err) => {
125 let completed_at = Utc::now();
126 let engine_err = EngineError::Operation(err);
127 if let Err(store_err) = self
128 .store
129 .update_step(
130 step.id,
131 StepUpdate {
132 status: Some(StepStatus::Failed),
133 error: Some(engine_err.to_string()),
134 completed_at: Some(completed_at),
135 ..StepUpdate::default()
136 },
137 )
138 .await
139 {
140 error!(step_id = %step.id, error = %store_err, "failed to persist step failure");
141 }
142
143 Err(engine_err)
144 }
145 }
146 }
147}