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//!
7//! A child that suspends (approval, human input, delay, signal) keeps its own
8//! suspension status and the parent's `Workflow` step stays open with the
9//! child run id in its output. The whole chain is suspended with it; when the
10//! root run replays, the open step re-enters the same child run.
11//!
12//! With `allow_failure` (see [`WorkflowContext::workflow_with`]) a child whose
13//! handler fails is still marked failed, but the parent's step completes with
14//! the failure in its [`SubWorkflowOutput`].
15
16use std::collections::HashMap;
17use std::time::Instant;
18
19use chrono::Utc;
20use rust_decimal::Decimal;
21use serde_json::{Value, from_value, json, to_value};
22use tracing::{error, info, warn};
23use uuid::Uuid;
24
25use ironflow_core::provider::LABEL_ROOT_RUN_ID;
26use ironflow_store::error::StoreError;
27use ironflow_store::models::{
28    NewRun, NewStep, RunStatus, RunUpdate, Step, StepKind, StepStatus, StepUpdate, TriggerKind,
29    step_trace_id,
30};
31
32use crate::config::{WorkflowOptions, WorkflowStepConfig};
33use crate::context::lifecycle::check_replay_identity;
34use crate::context::{PARENT_RUN_ID_LABEL, WorkflowContext};
35use crate::error::EngineError;
36use crate::executor::SubWorkflowOutput;
37use crate::guard::WorkflowRejection;
38use crate::handler::{TypedWorkflow, WorkflowHandler};
39use crate::plan::{SharedPlanRecorder, lock_plan};
40
41/// Key of the open `Workflow` step output that records the child run id.
42const CHILD_RUN_ID_KEY: &str = "child_run_id";
43
44/// The child run id recorded on an open `Workflow` step, if any.
45///
46/// A missing or unparsable id means the parent stopped before the child run
47/// was recorded: the step is reused but a new child run is started.
48fn recorded_child_run_id(step: &Step) -> Option<Uuid> {
49    let raw = step.output.as_ref()?.get(CHILD_RUN_ID_KEY)?.as_str()?;
50    match Uuid::parse_str(raw) {
51        Ok(id) => Some(id),
52        Err(err) => {
53            warn!(
54                step_id = %step.id,
55                value = %raw,
56                error = %err,
57                "open workflow step records an invalid child run id"
58            );
59            None
60        }
61    }
62}
63
64impl WorkflowContext {
65    /// Execute a sub-workflow step.
66    ///
67    /// Creates a child run of `handler` whose payload is `input`, executes it
68    /// with its own steps and lifecycle, and returns its run ID and aggregated
69    /// metrics. The child declares its input type through [`TypedWorkflow`],
70    /// so only a `W::Input` is accepted.
71    ///
72    /// When the child suspends (approval, human input, delay or signal), the
73    /// parent is suspended with it and this step stays open. Resuming the
74    /// child resumes the whole chain: the parent replays and re-enters the
75    /// same child run, whose completed steps are replayed.
76    ///
77    /// To tolerate a failed child, see [`workflow_with`](Self::workflow_with).
78    ///
79    /// Requires the context to be created with
80    /// `with_handler_resolver`.
81    ///
82    /// # Errors
83    ///
84    /// Returns [`EngineError::InvalidWorkflow`] if no handler is registered
85    /// with the given name, or if no handler resolver is available, and
86    /// [`EngineError::Serialization`] if `input` cannot be serialized. Returns
87    /// [`EngineError::ReplayDivergence`] when the step recorded at this
88    /// position has a different name or kind, and
89    /// [`EngineError::ChildSuspended`] when the child run suspended.
90    ///
91    /// # Examples
92    ///
93    /// ```no_run
94    /// use ironflow_engine::context::WorkflowContext;
95    /// use ironflow_engine::error::EngineError;
96    /// use ironflow_engine::handler::{HandlerFuture, TypedWorkflow, WorkflowHandler};
97    /// use serde::{Deserialize, Serialize};
98    ///
99    /// #[derive(Serialize, Deserialize)]
100    /// struct CollectInput {
101    ///     scope: String,
102    /// }
103    ///
104    /// struct Collect;
105    ///
106    /// impl WorkflowHandler for Collect {
107    ///     fn name(&self) -> &str { "collect" }
108    ///     fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
109    ///         Box::pin(async move { Ok(()) })
110    ///     }
111    /// }
112    ///
113    /// impl TypedWorkflow for Collect {
114    ///     type Input = CollectInput;
115    /// }
116    ///
117    /// # async fn example(ctx: &mut WorkflowContext) -> Result<(), EngineError> {
118    /// let child = ctx.workflow(&Collect, CollectInput { scope: "system".to_string() }).await?;
119    /// let steps = ctx.store().list_steps(child.run_id()).await?;
120    /// # Ok(())
121    /// # }
122    /// ```
123    ///
124    /// Any other input type is a compile error:
125    ///
126    /// ```compile_fail,E0308
127    /// # use ironflow_engine::context::WorkflowContext;
128    /// # use ironflow_engine::error::EngineError;
129    /// # use ironflow_engine::handler::{HandlerFuture, TypedWorkflow, WorkflowHandler};
130    /// # #[derive(serde::Serialize, serde::Deserialize)]
131    /// # struct CollectInput { scope: String }
132    /// # struct Collect;
133    /// # impl WorkflowHandler for Collect {
134    /// #     fn name(&self) -> &str { "collect" }
135    /// #     fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
136    /// #         Box::pin(async move { Ok(()) })
137    /// #     }
138    /// # }
139    /// # impl TypedWorkflow for Collect { type Input = CollectInput; }
140    /// # async fn example(ctx: &mut WorkflowContext) -> Result<(), EngineError> {
141    /// ctx.workflow(&Collect, serde_json::json!({"scope": "system"})).await?;
142    /// # Ok(())
143    /// # }
144    /// ```
145    pub async fn workflow<W: TypedWorkflow>(
146        &mut self,
147        handler: &W,
148        input: W::Input,
149    ) -> Result<SubWorkflowOutput, EngineError> {
150        let payload = to_value(&input)?;
151        self.run_sub_workflow(handler, payload, WorkflowOptions::default())
152            .await
153    }
154
155    /// Execute a sub-workflow step with [`WorkflowOptions`].
156    ///
157    /// Same as [`workflow`](Self::workflow). With
158    /// [`allow_failure`](WorkflowOptions::allow_failure), a child whose handler
159    /// fails does not fail the parent: the child run is still marked failed,
160    /// but this step completes with a [`SubWorkflowOutput`] whose
161    /// [`status`](SubWorkflowOutput::status) is `Failed` (or `Cancelled` when a
162    /// guardrail stopped it) and whose [`error`](SubWorkflowOutput::error)
163    /// carries the child error. The parent run then ends as `Warning`. A
164    /// resumed parent replays the completed step and creates no new child.
165    ///
166    /// A suspension is never tolerated: a child that suspends suspends the
167    /// parent. Errors raised outside the child (no resolver, unknown handler,
168    /// store errors, replay divergence, guard rejection of the invocation) are
169    /// not tolerated either.
170    ///
171    /// # Errors
172    ///
173    /// Same as [`workflow`](Self::workflow), except that the failure of the
174    /// child handler is returned in the output when `allow_failure` is set.
175    ///
176    /// # Examples
177    ///
178    /// ```no_run
179    /// use ironflow_engine::config::WorkflowOptions;
180    /// use ironflow_engine::context::WorkflowContext;
181    /// use ironflow_engine::error::EngineError;
182    /// use ironflow_engine::handler::{HandlerFuture, TypedWorkflow, WorkflowHandler};
183    /// use ironflow_store::entities::RunStatus;
184    /// use serde::{Deserialize, Serialize};
185    ///
186    /// #[derive(Serialize, Deserialize)]
187    /// struct CollectInput {
188    ///     scope: String,
189    /// }
190    ///
191    /// struct Collect;
192    ///
193    /// impl WorkflowHandler for Collect {
194    ///     fn name(&self) -> &str { "collect" }
195    ///     fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
196    ///         Box::pin(async move { Ok(()) })
197    ///     }
198    /// }
199    ///
200    /// impl TypedWorkflow for Collect {
201    ///     type Input = CollectInput;
202    /// }
203    ///
204    /// # async fn example(ctx: &mut WorkflowContext) -> Result<(), EngineError> {
205    /// let child = ctx
206    ///     .workflow_with(
207    ///         &Collect,
208    ///         CollectInput { scope: "system".to_string() },
209    ///         WorkflowOptions::new().allow_failure(),
210    ///     )
211    ///     .await?;
212    /// if child.status() == RunStatus::Failed {
213    ///     println!("collect failed: {:?}", child.error());
214    /// }
215    /// # Ok(())
216    /// # }
217    /// ```
218    pub async fn workflow_with<W: TypedWorkflow>(
219        &mut self,
220        handler: &W,
221        input: W::Input,
222        options: WorkflowOptions,
223    ) -> Result<SubWorkflowOutput, EngineError> {
224        let payload = to_value(&input)?;
225        self.run_sub_workflow(handler, payload, options).await
226    }
227
228    /// Execute a sub-workflow step whose child is only known at run time.
229    ///
230    /// Same as [`workflow`](Self::workflow), without the compile-time check of
231    /// the payload: the child must deserialize `payload` itself.
232    ///
233    /// # Errors
234    ///
235    /// Same as [`workflow`](Self::workflow).
236    ///
237    /// # Examples
238    ///
239    /// ```no_run
240    /// use ironflow_engine::context::WorkflowContext;
241    /// use ironflow_engine::error::EngineError;
242    /// use ironflow_engine::handler::WorkflowHandler;
243    /// use serde_json::json;
244    ///
245    /// # #[allow(deprecated)]
246    /// # async fn example(ctx: &mut WorkflowContext, child: &dyn WorkflowHandler) -> Result<(), EngineError> {
247    /// let result = ctx.workflow_dyn(child, json!({"scope": "system"})).await?;
248    /// println!("child run {}", result.run_id());
249    /// # Ok(())
250    /// # }
251    /// ```
252    #[deprecated(
253        note = "implement `TypedWorkflow` on the child and call `workflow`: its payload is then checked at compile time"
254    )]
255    pub async fn workflow_dyn(
256        &mut self,
257        handler: &dyn WorkflowHandler,
258        payload: Value,
259    ) -> Result<SubWorkflowOutput, EngineError> {
260        self.run_sub_workflow(handler, payload, WorkflowOptions::default())
261            .await
262    }
263
264    /// Record, then run or plan, a sub-workflow step.
265    ///
266    /// A `Workflow` step completed in a previous execution is replayed without
267    /// running the child again. A step left open (`Running`) by a suspended
268    /// child is reused and re-enters the child run it recorded.
269    async fn run_sub_workflow(
270        &mut self,
271        handler: &dyn WorkflowHandler,
272        payload: Value,
273        options: WorkflowOptions,
274    ) -> Result<SubWorkflowOutput, EngineError> {
275        // Plan mode: record the invocation, expand the child handler in the
276        // same recorder, and return a synthetic output. No child run is
277        // created and no step of the child is executed.
278        if let Some(plan) = self.plan().cloned() {
279            return self.plan_sub_workflow(&plan, handler, payload).await;
280        }
281
282        let mut config = WorkflowStepConfig::new(handler.name(), payload);
283        config.allow_failure = options.allow_failure;
284        let position = self.position;
285
286        let existing = self.replay_steps.get(&position).cloned();
287        if let Some(existing) = &existing {
288            check_replay_identity(
289                existing,
290                position,
291                &config.workflow_name,
292                &StepKind::Workflow,
293            )?;
294            if existing.status.state == StepStatus::Completed {
295                return self.replay_sub_workflow(existing);
296            }
297        }
298
299        // Guard check: verify limits before creating the step.
300        if let (Some(guard_config), Some(guard_state)) = (&self.guard_config, &self.guard_state) {
301            let state = guard_state
302                .lock()
303                .map_err(|_| WorkflowRejection::GuardUnavailable)?;
304            state.check(guard_config, handler.name())?;
305        }
306
307        self.position += 1;
308
309        // An open step was left by a child that suspended: reuse it instead of
310        // recording a second step at the same position.
311        let (step, resume) = match existing.filter(|s| s.status.state == StepStatus::Running) {
312            Some(step) => {
313                let resume = recorded_child_run_id(&step);
314                (step, resume)
315            }
316            None => {
317                let trace_id = step_trace_id(self.run_id, &config.workflow_name, position);
318                let step = self
319                    .store
320                    .create_step(NewStep {
321                        run_id: self.run_id,
322                        trace_id,
323                        name: config.workflow_name.clone(),
324                        kind: StepKind::Workflow,
325                        position,
326                        input: Some(to_value(&config)?),
327                        is_error_handler: false,
328                    })
329                    .await?;
330
331                self.start_step(step.id, Utc::now()).await?;
332                (step, None)
333            }
334        };
335
336        // Record invocation in guard state (fail-closed).
337        if let Some(guard_state) = &self.guard_state {
338            let mut state = guard_state
339                .lock()
340                .map_err(|_| WorkflowRejection::GuardUnavailable)?;
341            state.record_invocation(handler.name());
342        }
343
344        match self.execute_child_workflow(&config, step.id, resume).await {
345            Ok((output, child_had_allowed_failure)) => {
346                self.total_cost_usd += output.cost_usd();
347                self.total_duration_ms += output.duration_ms();
348                if child_had_allowed_failure {
349                    self.has_allowed_failure = true;
350                }
351
352                let completed_at = Utc::now();
353                self.store
354                    .update_step(
355                        step.id,
356                        StepUpdate {
357                            status: Some(StepStatus::Completed),
358                            output: Some(to_value(&output)?),
359                            duration_ms: Some(output.duration_ms()),
360                            cost_usd: Some(output.cost_usd()),
361                            completed_at: Some(completed_at),
362                            ..StepUpdate::default()
363                        },
364                    )
365                    .await?;
366
367                info!(
368                    run_id = %self.run_id,
369                    child_workflow = %config.workflow_name,
370                    duration_ms = output.duration_ms(),
371                    "workflow step completed"
372                );
373
374                self.last_step_ids = vec![step.id];
375
376                self.guard_record_return();
377                Ok(output)
378            }
379            // The child suspended: the step stays open, neither failed nor
380            // completed, so the next replay re-enters the same child run.
381            Err(err) if err.is_suspension() => {
382                self.guard_record_return();
383                Err(err)
384            }
385            Err(err) => {
386                let completed_at = Utc::now();
387                if let Err(store_err) = self
388                    .store
389                    .update_step(
390                        step.id,
391                        StepUpdate {
392                            status: Some(StepStatus::Failed),
393                            error: Some(err.to_string()),
394                            completed_at: Some(completed_at),
395                            ..StepUpdate::default()
396                        },
397                    )
398                    .await
399                {
400                    error!(step_id = %step.id, error = %store_err, "failed to persist step failure");
401                }
402
403                self.guard_record_return();
404                Err(err)
405            }
406        }
407    }
408
409    /// Replay a `Workflow` step completed in a previous execution: the child
410    /// run is not executed again and nothing is re-counted by the guard.
411    fn replay_sub_workflow(&mut self, step: &Step) -> Result<SubWorkflowOutput, EngineError> {
412        let recorded = step.output.clone().ok_or_else(|| {
413            EngineError::StepConfig(format!(
414                "completed workflow step {} has no recorded output",
415                step.id
416            ))
417        })?;
418        let output: SubWorkflowOutput = from_value(recorded)?;
419
420        self.position += 1;
421        // Cost is not added: `carry_over_run_totals` seeded `total_cost_usd`
422        // from the run totals persisted before the suspension, which already
423        // include this child.
424        self.total_duration_ms += output.duration_ms();
425        if matches!(
426            output.status(),
427            RunStatus::Warning | RunStatus::Failed | RunStatus::Cancelled
428        ) {
429            self.has_allowed_failure = true;
430        }
431        self.last_step_ids = vec![step.id];
432
433        info!(
434            run_id = %self.run_id,
435            child_run_id = %output.run_id(),
436            step = %step.name,
437            "workflow step replayed from previous execution"
438        );
439        Ok(output)
440    }
441
442    /// Record a sub-workflow invocation while planning, expanding the child
443    /// handler into the same plan when the depth limit allows it.
444    ///
445    /// The child plans against its own payload and under its own workflow
446    /// name; the parent's payload is restored on the way out.
447    async fn plan_sub_workflow(
448        &mut self,
449        plan: &SharedPlanRecorder,
450        handler: &dyn WorkflowHandler,
451        payload: Value,
452    ) -> Result<SubWorkflowOutput, EngineError> {
453        self.position += 1;
454        let sub_name = handler.name().to_string();
455        // No child run exists while planning: a nil id and zero metrics.
456        let planned = SubWorkflowOutput::new(
457            Uuid::nil(),
458            &sub_name,
459            RunStatus::Completed,
460            Decimal::ZERO,
461            0,
462        );
463
464        {
465            let mut recorder = lock_plan(plan);
466            if !recorder.record(&sub_name, StepKind::Workflow, &self.workflow_name, None) {
467                return Ok(planned);
468            }
469            recorder.set_last(vec![sub_name.clone()]);
470        }
471
472        let expand = lock_plan(plan).enter_workflow();
473        if expand {
474            let previous_payload = lock_plan(plan).swap_payload(payload.clone());
475
476            let mut child = WorkflowContext::new(
477                Uuid::now_v7(),
478                sub_name.clone(),
479                self.store.clone(),
480                self.provider.clone(),
481            );
482            child.handler_resolver = self.handler_resolver.clone();
483            child.set_plan(plan.clone());
484
485            if let Err(err) = handler.execute(&mut child).await {
486                lock_plan(plan).fail(format!(
487                    "sub-workflow {sub_name} could not be planned: {err}"
488                ));
489            }
490
491            let mut recorder = lock_plan(plan);
492            recorder.swap_payload(previous_payload);
493            recorder.leave_workflow();
494        }
495
496        Ok(planned)
497    }
498
499    /// Execute a child workflow and return aggregated output plus whether
500    /// at least one `allow_failure` step failed.
501    ///
502    /// `resume` is the child run recorded on an open step: that run is
503    /// re-entered, with its completed steps replayed, instead of creating a
504    /// new one. A child that suspends is left in its suspension status and
505    /// [`EngineError::ChildSuspended`] is returned.
506    async fn execute_child_workflow(
507        &self,
508        config: &WorkflowStepConfig,
509        step_id: Uuid,
510        resume: Option<Uuid>,
511    ) -> Result<(SubWorkflowOutput, bool), EngineError> {
512        let resolver = self.handler_resolver.as_ref().ok_or_else(|| {
513            EngineError::InvalidWorkflow(
514                "sub-workflow requires a handler resolver (use Engine to execute)".to_string(),
515            )
516        })?;
517
518        let handler = resolver(&config.workflow_name).ok_or_else(|| {
519            EngineError::InvalidWorkflow(format!("no handler registered: {}", config.workflow_name))
520        })?;
521
522        let (child_run_id, carried_cost_usd, carried_duration_ms) = match resume {
523            Some(child_run_id) => {
524                let child_run = self
525                    .store
526                    .get_run(child_run_id)
527                    .await?
528                    .ok_or(EngineError::Store(StoreError::RunNotFound(child_run_id)))?;
529
530                match child_run.status.state {
531                    // Already moved to Running by the path that resumed it.
532                    RunStatus::Running => {}
533                    RunStatus::AwaitingApproval | RunStatus::Pending => {
534                        self.store
535                            .update_run_status(child_run_id, RunStatus::Running)
536                            .await?;
537                    }
538                    RunStatus::Sleeping => {
539                        self.store
540                            .update_run_status(child_run_id, RunStatus::Pending)
541                            .await?;
542                        self.store
543                            .update_run_status(child_run_id, RunStatus::Running)
544                            .await?;
545                    }
546                    // The child finished but the parent stopped before closing
547                    // its step: report the recorded outcome, run nothing.
548                    status @ (RunStatus::Failed | RunStatus::Cancelled) if config.allow_failure => {
549                        let error = child_run
550                            .error
551                            .clone()
552                            .unwrap_or_else(|| "child run failed".to_string());
553                        return Ok((
554                            SubWorkflowOutput::new(
555                                child_run_id,
556                                &config.workflow_name,
557                                status,
558                                child_run.cost_usd,
559                                child_run.duration_ms,
560                            )
561                            .with_error(error),
562                            true,
563                        ));
564                    }
565                    status @ (RunStatus::Completed | RunStatus::Warning) => {
566                        return Ok((
567                            SubWorkflowOutput::new(
568                                child_run_id,
569                                &config.workflow_name,
570                                status,
571                                child_run.cost_usd,
572                                child_run.duration_ms,
573                            ),
574                            status == RunStatus::Warning,
575                        ));
576                    }
577                    other => {
578                        return Err(EngineError::InvalidWorkflow(format!(
579                            "child run {child_run_id} is {other}"
580                        )));
581                    }
582                }
583
584                info!(
585                    parent_run_id = %self.run_id,
586                    child_run_id = %child_run_id,
587                    workflow = %config.workflow_name,
588                    "child run re-entered"
589                );
590                (child_run_id, child_run.cost_usd, child_run.duration_ms)
591            }
592            None => {
593                let child_run_id = self.create_child_run(config).await?;
594
595                // Recorded before the child runs, so a suspension of the child
596                // can be resumed into this same run.
597                self.store
598                    .update_step(
599                        step_id,
600                        StepUpdate {
601                            output: Some(json!({ CHILD_RUN_ID_KEY: child_run_id })),
602                            ..StepUpdate::default()
603                        },
604                    )
605                    .await?;
606
607                self.store
608                    .update_run_status(child_run_id, RunStatus::Running)
609                    .await?;
610                (child_run_id, Decimal::ZERO, 0)
611            }
612        };
613
614        let run_start = Instant::now();
615        let mut child_ctx = WorkflowContext {
616            run_id: child_run_id,
617            root_run_id: self.root_run_id,
618            workflow_name: config.workflow_name.clone(),
619            store: self.store.clone(),
620            provider: self.provider.clone(),
621            decision_provider: self.decision_provider.clone(),
622            handler_resolver: self.handler_resolver.clone(),
623            position: 0,
624            last_step_ids: Vec::new(),
625            // A re-entered child starts from what it already spent, like a
626            // resumed top-level run.
627            total_cost_usd: carried_cost_usd,
628            total_duration_ms: 0,
629            max_cost_usd: self.max_cost_usd,
630            // Everything the parent chain already spent counts against the
631            // shared cap, so the child cannot restart the budget from zero.
632            inherited_cost_usd: self.charged_cost_usd(),
633            replay_steps: HashMap::new(),
634            replay_wave_steps: HashMap::new(),
635            granted_approvals: HashMap::new(),
636            answered_inputs: HashMap::new(),
637            // A child run is never itself retried.
638            attempt: 1,
639            carried_duration_ms,
640            log_sender: self.log_sender.clone(),
641            // A child shares the storage backend but not the parent's artifacts:
642            // input lookups are scoped to the child's own run.
643            artifact_sink: self.artifact_sink.clone(),
644            has_allowed_failure: false,
645            error_handlers: Vec::new(),
646            guard_state: self.guard_state.clone(),
647            guard_config: self.guard_config.clone(),
648            step_results: Vec::new(),
649            event_bus: self.event_bus.clone(),
650            // A child run is mocked exactly like its parent.
651            interceptor: self.interceptor.clone(),
652            trace_context: self.trace_context.child(),
653            operation_ctx: None,
654            run_created_at: None,
655            plan: None,
656        };
657
658        // A re-entered child replays its completed steps and is served the
659        // answer, signal or elapsed delay it was suspended on.
660        let loaded = if resume.is_some() {
661            child_ctx.load_replay_steps().await
662        } else {
663            Ok(())
664        };
665        let result = match loaded {
666            Ok(()) => handler.execute(&mut child_ctx).await,
667            Err(err) => Err(err),
668        };
669        let total_duration = child_ctx.carried_duration_ms + run_start.elapsed().as_millis() as u64;
670        let completed_at = Utc::now();
671
672        match result {
673            Ok(()) => {
674                let child_status = if child_ctx.has_allowed_failure {
675                    RunStatus::Warning
676                } else {
677                    RunStatus::Completed
678                };
679                self.store
680                    .update_run(
681                        child_run_id,
682                        RunUpdate {
683                            status: Some(child_status),
684                            cost_usd: Some(child_ctx.total_cost_usd),
685                            duration_ms: Some(total_duration),
686                            completed_at: Some(completed_at),
687                            ..RunUpdate::default()
688                        },
689                    )
690                    .await?;
691
692                let child_had_allowed_failure = child_ctx.has_allowed_failure;
693                Ok((
694                    SubWorkflowOutput::new(
695                        child_run_id,
696                        &config.workflow_name,
697                        child_status,
698                        child_ctx.total_cost_usd,
699                        total_duration,
700                    ),
701                    child_had_allowed_failure,
702                ))
703            }
704            Err(err) if err.is_suspension() => {
705                match self
706                    .suspend_child_run(child_run_id, &err, child_ctx.total_cost_usd, total_duration)
707                    .await
708                {
709                    Ok(()) => {
710                        info!(
711                            parent_run_id = %self.run_id,
712                            child_run_id = %child_run_id,
713                            cause = %err.suspension_leaf(),
714                            "child run suspended"
715                        );
716                        Err(EngineError::ChildSuspended {
717                            run_id: child_run_id,
718                            cause: Box::new(err),
719                        })
720                    }
721                    Err(store_err) => {
722                        self.fail_child_run(
723                            child_run_id,
724                            RunStatus::Failed,
725                            &store_err,
726                            child_ctx.total_cost_usd,
727                            total_duration,
728                        )
729                        .await;
730                        Err(store_err)
731                    }
732                }
733            }
734            Err(err) => {
735                // The engine cancels top-level runs stopped by a guardrail.
736                let status = if matches!(
737                    err,
738                    EngineError::RunBudgetExceeded { .. } | EngineError::WorkflowGuardRejected(_)
739                ) {
740                    RunStatus::Cancelled
741                } else {
742                    RunStatus::Failed
743                };
744                self.fail_child_run(
745                    child_run_id,
746                    status,
747                    &err,
748                    child_ctx.total_cost_usd,
749                    total_duration,
750                )
751                .await;
752                if config.allow_failure {
753                    return Ok((
754                        SubWorkflowOutput::new(
755                            child_run_id,
756                            &config.workflow_name,
757                            status,
758                            child_ctx.total_cost_usd,
759                            total_duration,
760                        )
761                        .with_error(err.to_string()),
762                        true,
763                    ));
764                }
765                Err(err)
766            }
767        }
768    }
769
770    /// Create the child run of a sub-workflow step and return its id.
771    ///
772    /// The child inherits the parent labels and author, and is linked to its
773    /// parent and to the root of the chain by two labels, so a suspended child
774    /// can be found and resumed like a top-level run.
775    async fn create_child_run(&self, config: &WorkflowStepConfig) -> Result<Uuid, EngineError> {
776        // Whoever triggered the parent workflow is accountable for its children.
777        let parent = self.store.get_run(self.run_id).await?;
778        let (mut labels, parent_author) =
779            parent.map(|r| (r.labels, r.created_by)).unwrap_or_default();
780        // Overwritten, never inherited: a grand-child must point at its own
781        // parent, not at its grand-parent.
782        labels.insert(PARENT_RUN_ID_LABEL.to_string(), self.run_id.to_string());
783        labels.insert(LABEL_ROOT_RUN_ID.to_string(), self.root_run_id.to_string());
784
785        let child_run = self
786            .store
787            .create_run(NewRun {
788                workflow_name: config.workflow_name.clone(),
789                trigger: TriggerKind::Workflow,
790                payload: config.payload.clone(),
791                max_retries: 0,
792                handler_version: None,
793                labels,
794                scheduled_at: None,
795                created_by: parent_author,
796                idempotency_key: None,
797                // The child shares the parent's cap; it does not get its own budget.
798                max_cost_usd: self.max_cost_usd,
799            })
800            .await?
801            .into_run();
802
803        info!(
804            parent_run_id = %self.run_id,
805            child_run_id = %child_run.id,
806            workflow = %config.workflow_name,
807            "child run created"
808        );
809        Ok(child_run.id)
810    }
811
812    /// Persist the suspension of a child run, with no event: the root run
813    /// publishes the suspension once the whole chain is suspended.
814    ///
815    /// A direct suspension is persisted like a top-level run's (a delay or a
816    /// signal deadline arms `scheduled_at`). A child suspended because of its
817    /// own child gets no `scheduled_at`: only the deepest run owns the
818    /// wake-up, so the chain is never resumed twice.
819    async fn suspend_child_run(
820        &self,
821        child_run_id: Uuid,
822        err: &EngineError,
823        cost_usd: Decimal,
824        duration_ms: u64,
825    ) -> Result<(), EngineError> {
826        let totals = RunUpdate {
827            cost_usd: Some(cost_usd),
828            duration_ms: Some(duration_ms),
829            ..RunUpdate::default()
830        };
831
832        let update = match err {
833            EngineError::DelaySleeping { wake_at, .. } => RunUpdate {
834                status: Some(RunStatus::Sleeping),
835                scheduled_at: Some(*wake_at),
836                ..totals
837            },
838            EngineError::SignalWaiting {
839                step_id,
840                deadline_at,
841                ..
842            } => {
843                // Atomic with the step lock, like a top-level run.
844                self.store
845                    .suspend_run_on_signal(child_run_id, *step_id, *deadline_at)
846                    .await?;
847                totals
848            }
849            EngineError::ChildSuspended { cause, .. } => RunUpdate {
850                status: Some(cause.suspension_status()),
851                ..totals
852            },
853            _ => RunUpdate {
854                status: Some(RunStatus::AwaitingApproval),
855                ..totals
856            },
857        };
858
859        self.store.update_run(child_run_id, update).await?;
860        Ok(())
861    }
862
863    /// Mark a child run failed after its handler (or its suspension) failed.
864    ///
865    /// Best effort: the original error is what the parent reports.
866    async fn fail_child_run(
867        &self,
868        child_run_id: Uuid,
869        status: RunStatus,
870        err: &EngineError,
871        cost_usd: Decimal,
872        duration_ms: u64,
873    ) {
874        if let Err(store_err) = self
875            .store
876            .update_run(
877                child_run_id,
878                RunUpdate {
879                    status: Some(status),
880                    error: Some(err.to_string()),
881                    cost_usd: Some(cost_usd),
882                    duration_ms: Some(duration_ms),
883                    completed_at: Some(Utc::now()),
884                    ..RunUpdate::default()
885                },
886            )
887            .await
888        {
889            error!(
890                child_run_id = %child_run_id,
891                store_error = %store_err,
892                "failed to persist child run failure"
893            );
894        }
895    }
896}