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