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