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;
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}