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::{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}