Skip to main content

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}