Skip to main content

ironflow_engine/context/steps/
sub_workflow.rs

1//! Sub-workflow step for [`WorkflowContext`].
2//!
3//! A sub-workflow runs a registered [`WorkflowHandler`] in its own child run.
4//! The child context is built here from the parent's private fields, which is
5//! possible because this module is a descendant of `context`.
6
7use std::collections::HashMap;
8use std::time::Instant;
9
10use chrono::Utc;
11use rust_decimal::Decimal;
12use serde_json::{Value, to_value};
13use tracing::{error, info};
14use uuid::Uuid;
15
16use ironflow_store::models::{
17    NewRun, NewStep, RunStatus, RunUpdate, StepKind, StepStatus, StepUpdate, TriggerKind,
18    step_trace_id,
19};
20
21use crate::config::WorkflowStepConfig;
22use crate::context::WorkflowContext;
23use crate::error::EngineError;
24use crate::executor::SubWorkflowOutput;
25use crate::guard::WorkflowRejection;
26use crate::handler::{TypedWorkflow, WorkflowHandler};
27use crate::plan::{SharedPlanRecorder, lock_plan};
28
29impl WorkflowContext {
30    /// Execute a sub-workflow step.
31    ///
32    /// Creates a child run of `handler` whose payload is `input`, executes it
33    /// with its own steps and lifecycle, and returns its run ID and aggregated
34    /// metrics. The child declares its input type through [`TypedWorkflow`],
35    /// so only a `W::Input` is accepted.
36    ///
37    /// Requires the context to be created with
38    /// `with_handler_resolver`.
39    ///
40    /// # Errors
41    ///
42    /// Returns [`EngineError::InvalidWorkflow`] if no handler is registered
43    /// with the given name, or if no handler resolver is available, and
44    /// [`EngineError::Serialization`] if `input` cannot be serialized.
45    ///
46    /// # Examples
47    ///
48    /// ```no_run
49    /// use ironflow_engine::context::WorkflowContext;
50    /// use ironflow_engine::error::EngineError;
51    /// use ironflow_engine::handler::{HandlerFuture, TypedWorkflow, WorkflowHandler};
52    /// use serde::{Deserialize, Serialize};
53    ///
54    /// #[derive(Serialize, Deserialize)]
55    /// struct CollectInput {
56    ///     scope: String,
57    /// }
58    ///
59    /// struct Collect;
60    ///
61    /// impl WorkflowHandler for Collect {
62    ///     fn name(&self) -> &str { "collect" }
63    ///     fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
64    ///         Box::pin(async move { Ok(()) })
65    ///     }
66    /// }
67    ///
68    /// impl TypedWorkflow for Collect {
69    ///     type Input = CollectInput;
70    /// }
71    ///
72    /// # async fn example(ctx: &mut WorkflowContext) -> Result<(), EngineError> {
73    /// let child = ctx.workflow(&Collect, CollectInput { scope: "system".to_string() }).await?;
74    /// let steps = ctx.store().list_steps(child.run_id()).await?;
75    /// # Ok(())
76    /// # }
77    /// ```
78    ///
79    /// Any other input type is a compile error:
80    ///
81    /// ```compile_fail,E0308
82    /// # use ironflow_engine::context::WorkflowContext;
83    /// # use ironflow_engine::error::EngineError;
84    /// # use ironflow_engine::handler::{HandlerFuture, TypedWorkflow, WorkflowHandler};
85    /// # #[derive(serde::Serialize, serde::Deserialize)]
86    /// # struct CollectInput { scope: String }
87    /// # struct Collect;
88    /// # impl WorkflowHandler for Collect {
89    /// #     fn name(&self) -> &str { "collect" }
90    /// #     fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
91    /// #         Box::pin(async move { Ok(()) })
92    /// #     }
93    /// # }
94    /// # impl TypedWorkflow for Collect { type Input = CollectInput; }
95    /// # async fn example(ctx: &mut WorkflowContext) -> Result<(), EngineError> {
96    /// ctx.workflow(&Collect, serde_json::json!({"scope": "system"})).await?;
97    /// # Ok(())
98    /// # }
99    /// ```
100    pub async fn workflow<W: TypedWorkflow>(
101        &mut self,
102        handler: &W,
103        input: W::Input,
104    ) -> Result<SubWorkflowOutput, EngineError> {
105        let payload = to_value(&input)?;
106        self.run_sub_workflow(handler, payload).await
107    }
108
109    /// Execute a sub-workflow step whose child is only known at run time.
110    ///
111    /// Same as [`workflow`](Self::workflow), without the compile-time check of
112    /// the payload: the child must deserialize `payload` itself.
113    ///
114    /// # Errors
115    ///
116    /// Same as [`workflow`](Self::workflow).
117    ///
118    /// # Examples
119    ///
120    /// ```no_run
121    /// use ironflow_engine::context::WorkflowContext;
122    /// use ironflow_engine::error::EngineError;
123    /// use ironflow_engine::handler::WorkflowHandler;
124    /// use serde_json::json;
125    ///
126    /// # #[allow(deprecated)]
127    /// # async fn example(ctx: &mut WorkflowContext, child: &dyn WorkflowHandler) -> Result<(), EngineError> {
128    /// let result = ctx.workflow_dyn(child, json!({"scope": "system"})).await?;
129    /// println!("child run {}", result.run_id());
130    /// # Ok(())
131    /// # }
132    /// ```
133    #[deprecated(
134        note = "implement `TypedWorkflow` on the child and call `workflow`: its payload is then checked at compile time"
135    )]
136    pub async fn workflow_dyn(
137        &mut self,
138        handler: &dyn WorkflowHandler,
139        payload: Value,
140    ) -> Result<SubWorkflowOutput, EngineError> {
141        self.run_sub_workflow(handler, payload).await
142    }
143
144    /// Record, then run or plan, a sub-workflow step.
145    async fn run_sub_workflow(
146        &mut self,
147        handler: &dyn WorkflowHandler,
148        payload: Value,
149    ) -> Result<SubWorkflowOutput, EngineError> {
150        // Plan mode: record the invocation, expand the child handler in the
151        // same recorder, and return a synthetic output. No child run is
152        // created and no step of the child is executed.
153        if let Some(plan) = self.plan().cloned() {
154            return self.plan_sub_workflow(&plan, handler, payload).await;
155        }
156
157        // Guard check: verify limits before creating the step.
158        if let (Some(guard_config), Some(guard_state)) = (&self.guard_config, &self.guard_state) {
159            let state = guard_state
160                .lock()
161                .map_err(|_| WorkflowRejection::GuardUnavailable)?;
162            state.check(guard_config, handler.name())?;
163        }
164
165        let config = WorkflowStepConfig::new(handler.name(), payload);
166        let position = self.position;
167        self.position += 1;
168
169        let trace_id = step_trace_id(self.run_id, &config.workflow_name, position);
170        let step = self
171            .store
172            .create_step(NewStep {
173                run_id: self.run_id,
174                trace_id,
175                name: config.workflow_name.clone(),
176                kind: StepKind::Workflow,
177                position,
178                input: Some(to_value(&config)?),
179                is_error_handler: false,
180            })
181            .await?;
182
183        self.start_step(step.id, Utc::now()).await?;
184
185        // Record invocation in guard state (fail-closed).
186        if let Some(guard_state) = &self.guard_state {
187            let mut state = guard_state
188                .lock()
189                .map_err(|_| WorkflowRejection::GuardUnavailable)?;
190            state.record_invocation(handler.name());
191        }
192
193        match self.execute_child_workflow(&config).await {
194            Ok((output, child_had_allowed_failure)) => {
195                self.total_cost_usd += output.cost_usd();
196                self.total_duration_ms += output.duration_ms();
197                if child_had_allowed_failure {
198                    self.has_allowed_failure = true;
199                }
200
201                let completed_at = Utc::now();
202                self.store
203                    .update_step(
204                        step.id,
205                        StepUpdate {
206                            status: Some(StepStatus::Completed),
207                            output: Some(to_value(&output)?),
208                            duration_ms: Some(output.duration_ms()),
209                            cost_usd: Some(output.cost_usd()),
210                            completed_at: Some(completed_at),
211                            ..StepUpdate::default()
212                        },
213                    )
214                    .await?;
215
216                info!(
217                    run_id = %self.run_id,
218                    child_workflow = %config.workflow_name,
219                    duration_ms = output.duration_ms(),
220                    "workflow step completed"
221                );
222
223                self.last_step_ids = vec![step.id];
224
225                self.guard_record_return();
226                Ok(output)
227            }
228            Err(err) => {
229                let completed_at = Utc::now();
230                if let Err(store_err) = self
231                    .store
232                    .update_step(
233                        step.id,
234                        StepUpdate {
235                            status: Some(StepStatus::Failed),
236                            error: Some(err.to_string()),
237                            completed_at: Some(completed_at),
238                            ..StepUpdate::default()
239                        },
240                    )
241                    .await
242                {
243                    error!(step_id = %step.id, error = %store_err, "failed to persist step failure");
244                }
245
246                self.guard_record_return();
247                Err(err)
248            }
249        }
250    }
251
252    /// Record a sub-workflow invocation while planning, expanding the child
253    /// handler into the same plan when the depth limit allows it.
254    ///
255    /// The child plans against its own payload and under its own workflow
256    /// name; the parent's payload is restored on the way out.
257    async fn plan_sub_workflow(
258        &mut self,
259        plan: &SharedPlanRecorder,
260        handler: &dyn WorkflowHandler,
261        payload: Value,
262    ) -> Result<SubWorkflowOutput, EngineError> {
263        self.position += 1;
264        let sub_name = handler.name().to_string();
265        // No child run exists while planning: a nil id and zero metrics.
266        let planned = SubWorkflowOutput::new(
267            Uuid::nil(),
268            &sub_name,
269            RunStatus::Completed,
270            Decimal::ZERO,
271            0,
272        );
273
274        {
275            let mut recorder = lock_plan(plan);
276            if !recorder.record(&sub_name, StepKind::Workflow, &self.workflow_name, None) {
277                return Ok(planned);
278            }
279            recorder.set_last(vec![sub_name.clone()]);
280        }
281
282        let expand = lock_plan(plan).enter_workflow();
283        if expand {
284            let previous_payload = lock_plan(plan).swap_payload(payload.clone());
285
286            let mut child = WorkflowContext::new(
287                Uuid::now_v7(),
288                sub_name.clone(),
289                self.store.clone(),
290                self.provider.clone(),
291            );
292            child.handler_resolver = self.handler_resolver.clone();
293            child.set_plan(plan.clone());
294
295            if let Err(err) = handler.execute(&mut child).await {
296                lock_plan(plan).fail(format!(
297                    "sub-workflow {sub_name} could not be planned: {err}"
298                ));
299            }
300
301            let mut recorder = lock_plan(plan);
302            recorder.swap_payload(previous_payload);
303            recorder.leave_workflow();
304        }
305
306        Ok(planned)
307    }
308
309    /// Execute a child workflow and return aggregated output plus whether
310    /// at least one `allow_failure` step failed.
311    async fn execute_child_workflow(
312        &self,
313        config: &WorkflowStepConfig,
314    ) -> Result<(SubWorkflowOutput, bool), EngineError> {
315        let resolver = self.handler_resolver.as_ref().ok_or_else(|| {
316            EngineError::InvalidWorkflow(
317                "sub-workflow requires a handler resolver (use Engine to execute)".to_string(),
318            )
319        })?;
320
321        let handler = resolver(&config.workflow_name).ok_or_else(|| {
322            EngineError::InvalidWorkflow(format!("no handler registered: {}", config.workflow_name))
323        })?;
324
325        // A child run inherits both the parent labels and the parent author:
326        // whoever triggered the parent workflow is accountable for its children.
327        let parent = self.store.get_run(self.run_id).await?;
328        let (parent_labels, parent_author) =
329            parent.map(|r| (r.labels, r.created_by)).unwrap_or_default();
330
331        let child_run = self
332            .store
333            .create_run(NewRun {
334                workflow_name: config.workflow_name.clone(),
335                trigger: TriggerKind::Workflow,
336                payload: config.payload.clone(),
337                max_retries: 0,
338                handler_version: None,
339                labels: parent_labels,
340                scheduled_at: None,
341                created_by: parent_author,
342                idempotency_key: None,
343                // The child shares the parent's cap; it does not get its own budget.
344                max_cost_usd: self.max_cost_usd,
345            })
346            .await?
347            .into_run();
348
349        let child_run_id = child_run.id;
350        info!(
351            parent_run_id = %self.run_id,
352            child_run_id = %child_run_id,
353            workflow = %config.workflow_name,
354            "child run created"
355        );
356
357        self.store
358            .update_run_status(child_run_id, RunStatus::Running)
359            .await?;
360
361        let run_start = Instant::now();
362        let mut child_ctx = WorkflowContext {
363            run_id: child_run_id,
364            workflow_name: config.workflow_name.clone(),
365            store: self.store.clone(),
366            provider: self.provider.clone(),
367            decision_provider: self.decision_provider.clone(),
368            handler_resolver: self.handler_resolver.clone(),
369            position: 0,
370            last_step_ids: Vec::new(),
371            total_cost_usd: Decimal::ZERO,
372            total_duration_ms: 0,
373            max_cost_usd: self.max_cost_usd,
374            // Everything the parent chain already spent counts against the
375            // shared cap, so the child cannot restart the budget from zero.
376            inherited_cost_usd: self.charged_cost_usd(),
377            replay_steps: HashMap::new(),
378            granted_approvals: HashMap::new(),
379            // A child run is created fresh here; it is never itself retried.
380            attempt: 1,
381            carried_duration_ms: 0,
382            log_sender: self.log_sender.clone(),
383            // A child shares the storage backend but not the parent's artifacts:
384            // input lookups are scoped to the child's own run.
385            artifact_sink: self.artifact_sink.clone(),
386            has_allowed_failure: false,
387            error_handlers: Vec::new(),
388            guard_state: self.guard_state.clone(),
389            guard_config: self.guard_config.clone(),
390            step_results: Vec::new(),
391            event_bus: self.event_bus.clone(),
392            // A child run is mocked exactly like its parent.
393            interceptor: self.interceptor.clone(),
394            trace_context: self.trace_context.child(),
395            operation_ctx: None,
396            plan: None,
397        };
398
399        let result = handler.execute(&mut child_ctx).await;
400        let total_duration = run_start.elapsed().as_millis() as u64;
401        let completed_at = Utc::now();
402
403        match result {
404            Ok(()) => {
405                let child_status = if child_ctx.has_allowed_failure {
406                    RunStatus::Warning
407                } else {
408                    RunStatus::Completed
409                };
410                self.store
411                    .update_run(
412                        child_run_id,
413                        RunUpdate {
414                            status: Some(child_status),
415                            cost_usd: Some(child_ctx.total_cost_usd),
416                            duration_ms: Some(total_duration),
417                            completed_at: Some(completed_at),
418                            ..RunUpdate::default()
419                        },
420                    )
421                    .await?;
422
423                let child_had_allowed_failure = child_ctx.has_allowed_failure;
424                Ok((
425                    SubWorkflowOutput::new(
426                        child_run_id,
427                        &config.workflow_name,
428                        child_status,
429                        child_ctx.total_cost_usd,
430                        total_duration,
431                    ),
432                    child_had_allowed_failure,
433                ))
434            }
435            Err(err) => {
436                if let Err(store_err) = self
437                    .store
438                    .update_run(
439                        child_run_id,
440                        RunUpdate {
441                            status: Some(RunStatus::Failed),
442                            error: Some(err.to_string()),
443                            cost_usd: Some(child_ctx.total_cost_usd),
444                            duration_ms: Some(total_duration),
445                            completed_at: Some(completed_at),
446                            ..RunUpdate::default()
447                        },
448                    )
449                    .await
450                {
451                    error!(
452                        child_run_id = %child_run_id,
453                        store_error = %store_err,
454                        "failed to persist child run failure"
455                    );
456                }
457
458                Err(err)
459            }
460        }
461    }
462}