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