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    /// On resume, a step that already completed in a prior execution of the
28    /// current attempt is replayed from the store instead of calling
29    /// [`Operation::execute`] again.
30    ///
31    /// # Errors
32    ///
33    /// Returns [`EngineError`] if the operation fails or the store errors.
34    /// Returns [`EngineError::ReplayDivergence`] when the step recorded at
35    /// this position has a different name or kind.
36    ///
37    /// # Examples
38    ///
39    /// ```no_run
40    /// use async_trait::async_trait;
41    /// use ironflow_engine::context::WorkflowContext;
42    /// use ironflow_engine::operation::{Operation, OperationContext};
43    /// use ironflow_core::error::OperationError;
44    /// use ironflow_engine::error::EngineError;
45    /// use serde_json::{Value, json};
46    ///
47    /// struct MyOp;
48    /// #[async_trait]
49    /// impl Operation for MyOp {
50    ///     fn kind(&self) -> &str { "my-service" }
51    ///     async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
52    ///         Ok(json!({"ok": true}))
53    ///     }
54    /// }
55    ///
56    /// # async fn example(ctx: &mut WorkflowContext) -> Result<(), EngineError> {
57    /// let result = ctx.operation("call-service", &MyOp).await?;
58    /// println!("output: {}", result.output);
59    /// # Ok(())
60    /// # }
61    /// ```
62    pub async fn operation(
63        &mut self,
64        name: &str,
65        op: &dyn Operation,
66    ) -> Result<StepOutput, EngineError> {
67        let kind = StepKind::Custom(op.kind().to_string());
68
69        // Plan mode: `op.execute` is never called, so no third-party API is
70        // touched while planning.
71        if let Some(plan) = self.plan().cloned() {
72            self.position += 1;
73            let mut recorder = lock_plan(&plan);
74            let estimate = recorder.estimate_for(name);
75            if recorder.record(name, kind, &self.workflow_name, None) {
76                recorder.set_last(vec![name.to_string()]);
77            }
78            let mut output = planned_custom_output(estimate);
79            output.artifacts = StepArtifacts::new(name, None, &[]);
80            return Ok(output);
81        }
82
83        let position = self.position;
84        self.position += 1;
85
86        // Replay: if this step already completed in a prior execution of the
87        // current attempt, return its cached output without calling
88        // `op.execute` or creating a new step.
89        if let Some(mut output) = self.try_replay_step(position, name, &kind)? {
90            let step_id = self.last_step_ids.last().copied();
91            output.artifacts = StepArtifacts::new(name, step_id, &[]);
92            return Ok(output);
93        }
94
95        let trace_id = step_trace_id(self.run_id, name, position);
96        let step = self
97            .store
98            .create_step(NewStep {
99                run_id: self.run_id,
100                trace_id,
101                name: name.to_string(),
102                kind,
103                position,
104                input: op.input(),
105                is_error_handler: false,
106            })
107            .await?;
108
109        self.start_step(step.id, Utc::now()).await?;
110
111        let start = Instant::now();
112
113        let op_ctx = self.ensure_operation_ctx();
114
115        match op.execute(op_ctx).await {
116            Ok(output_value) => {
117                let duration_ms = start.elapsed().as_millis() as u64;
118                self.total_duration_ms += duration_ms;
119
120                let completed_at = Utc::now();
121                self.store
122                    .update_step(
123                        step.id,
124                        StepUpdate {
125                            status: Some(StepStatus::Completed),
126                            output: Some(output_value.clone()),
127                            duration_ms: Some(duration_ms),
128                            cost_usd: Some(Decimal::ZERO),
129                            completed_at: Some(completed_at),
130                            ..StepUpdate::default()
131                        },
132                    )
133                    .await?;
134
135                info!(
136                    run_id = %self.run_id,
137                    step = %name,
138                    kind = op.kind(),
139                    duration_ms,
140                    "operation step completed"
141                );
142
143                self.last_step_ids = vec![step.id];
144
145                Ok(StepOutput {
146                    output: output_value,
147                    duration_ms,
148                    cost_usd: Decimal::ZERO,
149                    input_tokens: None,
150                    cache_read_input_tokens: None,
151                    cache_creation_input_tokens: None,
152                    output_tokens: None,
153                    model: None,
154                    debug_messages: None,
155                    // An operation declares no output; its record is what
156                    // `put_artifact` attaches bytes to.
157                    artifacts: StepArtifacts::new(name, Some(step.id), &[]),
158                })
159            }
160            Err(err) => {
161                let completed_at = Utc::now();
162                let engine_err = EngineError::Operation(err);
163                if let Err(store_err) = self
164                    .store
165                    .update_step(
166                        step.id,
167                        StepUpdate {
168                            status: Some(StepStatus::Failed),
169                            error: Some(engine_err.to_string()),
170                            completed_at: Some(completed_at),
171                            ..StepUpdate::default()
172                        },
173                    )
174                    .await
175                {
176                    error!(step_id = %step.id, error = %store_err, "failed to persist step failure");
177                }
178
179                Err(engine_err)
180            }
181        }
182    }
183}