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, clamp_priority};
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        let (concurrency_key, priority) = options.into_parts();
379        config.concurrency_key = concurrency_key;
380        config.priority = priority;
381        let position = self.position;
382
383        let existing = self.replay_steps.get(&position).cloned();
384        if let Some(existing) = &existing {
385            check_replay_identity(
386                existing,
387                position,
388                &config.workflow_name,
389                &StepKind::Workflow,
390            )?;
391            if existing.status.state == StepStatus::Completed {
392                return self.replay_sub_workflow(existing);
393            }
394        }
395
396        // A step interrupted by a lost lease must be the same step the handler
397        // calls now before its child run is re-entered.
398        let interrupted = self.interrupted_children.get(&position).cloned();
399        if let Some(interrupted) = &interrupted {
400            check_replay_identity(
401                interrupted,
402                position,
403                &config.workflow_name,
404                &StepKind::Workflow,
405            )?;
406        }
407
408        // A child runs on the worker that runs its parent: refuse it here
409        // rather than run it on a host missing what it requires.
410        self.check_child_worker_tags(handler)?;
411
412        // Guard check: verify limits before creating the step.
413        if let (Some(guard_config), Some(guard_state)) = (&self.guard_config, &self.guard_state) {
414            let state = guard_state
415                .lock()
416                .map_err(|_| WorkflowRejection::GuardUnavailable)?;
417            state.check(guard_config, handler.name())?;
418        }
419
420        // Paused while it ran: no child is started, the run stops here.
421        self.check_not_paused().await?;
422
423        self.position += 1;
424
425        // An open step was left by a child that suspended: reuse it instead of
426        // recording a second step at the same position.
427        let (step, resume) = match existing.filter(|s| s.status.state == StepStatus::Running) {
428            Some(step) => {
429                let resume = recorded_child_run_id(&step).map(ChildResume::Suspended);
430                (step, resume)
431            }
432            None => {
433                let trace_id = step_trace_id(self.run_id, &config.workflow_name, position);
434                let step = self
435                    .store
436                    .create_step(NewStep {
437                        run_id: self.run_id,
438                        trace_id,
439                        name: config.workflow_name.clone(),
440                        kind: StepKind::Workflow,
441                        position,
442                        input: Some(to_value(&config)?),
443                        is_error_handler: false,
444                    })
445                    .await?;
446
447                self.start_step(step.id, Utc::now()).await?;
448                // A step interrupted by a lost lease re-enters the child run it
449                // recorded instead of starting a new one.
450                let resume = interrupted
451                    .as_ref()
452                    .and_then(recorded_child_run_id)
453                    .map(ChildResume::Interrupted);
454                (step, resume)
455            }
456        };
457
458        // Record invocation in guard state (fail-closed).
459        if let Some(guard_state) = &self.guard_state {
460            let mut state = guard_state
461                .lock()
462                .map_err(|_| WorkflowRejection::GuardUnavailable)?;
463            state.record_invocation(handler.name());
464        }
465
466        match self.execute_child_workflow(&config, step.id, resume).await {
467            // No child run was created: the step completes with the conflict
468            // as its output, so a replay serves the same outcome.
469            Ok(ChildOutcome::Conflict(conflict)) => {
470                self.store
471                    .update_step(
472                        step.id,
473                        StepUpdate {
474                            status: Some(StepStatus::Completed),
475                            output: Some(json!({ "concurrency_conflict": conflict })),
476                            duration_ms: Some(0),
477                            cost_usd: Some(Decimal::ZERO),
478                            completed_at: Some(Utc::now()),
479                            ..StepUpdate::default()
480                        },
481                    )
482                    .await?;
483
484                info!(
485                    run_id = %self.run_id,
486                    child_workflow = %config.workflow_name,
487                    key = %conflict.key(),
488                    holder = %conflict.run_id(),
489                    "workflow step skipped: concurrency conflict"
490                );
491
492                self.last_step_ids = vec![step.id];
493
494                self.guard_record_return();
495                Ok(SubWorkflowOutcome::Conflict(conflict))
496            }
497            Ok(ChildOutcome::Finished(output, child_had_allowed_failure)) => {
498                self.total_cost_usd += output.cost_usd();
499                self.total_duration_ms += output.duration_ms();
500                if child_had_allowed_failure {
501                    self.has_allowed_failure = true;
502                }
503
504                let completed_at = Utc::now();
505                self.store
506                    .update_step(
507                        step.id,
508                        StepUpdate {
509                            status: Some(StepStatus::Completed),
510                            output: Some(to_value(&output)?),
511                            duration_ms: Some(output.duration_ms()),
512                            cost_usd: Some(output.cost_usd()),
513                            completed_at: Some(completed_at),
514                            ..StepUpdate::default()
515                        },
516                    )
517                    .await?;
518
519                info!(
520                    run_id = %self.run_id,
521                    child_workflow = %config.workflow_name,
522                    duration_ms = output.duration_ms(),
523                    "workflow step completed"
524                );
525
526                self.last_step_ids = vec![step.id];
527
528                self.guard_record_return();
529                Ok(SubWorkflowOutcome::Completed(output))
530            }
531            // The child suspended or was paused: the step stays open, neither
532            // failed nor completed, so the next replay re-enters the same child run.
533            Err(err) if err.is_suspension() || matches!(err, EngineError::RunPaused { .. }) => {
534                self.guard_record_return();
535                Err(err)
536            }
537            Err(err) => {
538                let completed_at = Utc::now();
539                if let Err(store_err) = self
540                    .store
541                    .update_step(
542                        step.id,
543                        StepUpdate {
544                            status: Some(StepStatus::Failed),
545                            error: Some(err.to_string()),
546                            completed_at: Some(completed_at),
547                            ..StepUpdate::default()
548                        },
549                    )
550                    .await
551                {
552                    error!(step_id = %step.id, error = %store_err, "failed to persist step failure");
553                }
554
555                self.guard_record_return();
556                Err(err)
557            }
558        }
559    }
560
561    /// Replay a `Workflow` step completed in a previous execution: the child
562    /// run is not executed again and nothing is re-counted by the guard.
563    ///
564    /// A step skipped on a concurrency conflict replays the same conflict: the
565    /// child is not attempted again, even if the key has been released since.
566    fn replay_sub_workflow(&mut self, step: &Step) -> Result<SubWorkflowOutcome, EngineError> {
567        let recorded = step.output.clone().ok_or_else(|| {
568            EngineError::StepConfig(format!(
569                "completed workflow step {} has no recorded output",
570                step.id
571            ))
572        })?;
573        let recorded: RecordedWorkflowStep = from_value(recorded)?;
574
575        self.position += 1;
576        self.last_step_ids = vec![step.id];
577
578        let output = match SubWorkflowOutcome::from(recorded) {
579            SubWorkflowOutcome::Completed(output) => output,
580            SubWorkflowOutcome::Conflict(conflict) => {
581                info!(
582                    run_id = %self.run_id,
583                    step = %step.name,
584                    key = %conflict.key(),
585                    holder = %conflict.run_id(),
586                    "workflow step replayed: concurrency conflict"
587                );
588                return Ok(SubWorkflowOutcome::Conflict(conflict));
589            }
590        };
591
592        // Cost is not added: `carry_over_run_totals` seeded `total_cost_usd`
593        // from the run totals persisted before the suspension, which already
594        // include this child.
595        self.total_duration_ms += output.duration_ms();
596        if matches!(
597            output.status(),
598            RunStatus::Warning | RunStatus::Failed | RunStatus::Cancelled
599        ) {
600            self.has_allowed_failure = true;
601        }
602
603        info!(
604            run_id = %self.run_id,
605            child_run_id = %output.run_id(),
606            step = %step.name,
607            "workflow step replayed from previous execution"
608        );
609        Ok(SubWorkflowOutcome::Completed(output))
610    }
611
612    /// Refuse a child whose required worker tags are not all carried by the
613    /// worker running this context.
614    ///
615    /// No check is made outside a tagged worker (API or local mode).
616    fn check_child_worker_tags(&self, handler: &dyn WorkflowHandler) -> Result<(), EngineError> {
617        let Some(carried) = &self.worker_tags else {
618            return Ok(());
619        };
620        let missing: Vec<String> = normalize_worker_tags(handler.required_worker_tags())
621            .into_iter()
622            .filter(|tag| !carried.contains(tag))
623            .collect();
624        if missing.is_empty() {
625            return Ok(());
626        }
627        Err(EngineError::InvalidWorkflow(format!(
628            "sub-workflow '{}' requires worker tags [{}] that this worker does not carry",
629            handler.name(),
630            missing.join(", ")
631        )))
632    }
633
634    /// Record a sub-workflow invocation while planning, expanding the child
635    /// handler into the same plan when the depth limit allows it.
636    ///
637    /// The child plans against its own payload and under its own workflow
638    /// name; the parent's payload is restored on the way out.
639    async fn plan_sub_workflow(
640        &mut self,
641        plan: &SharedPlanRecorder,
642        handler: &dyn WorkflowHandler,
643        payload: Value,
644    ) -> Result<SubWorkflowOutput, EngineError> {
645        self.position += 1;
646        let sub_name = handler.name().to_string();
647        // No child run exists while planning: a nil id and zero metrics.
648        let planned = SubWorkflowOutput::new(
649            Uuid::nil(),
650            &sub_name,
651            RunStatus::Completed,
652            Decimal::ZERO,
653            0,
654        );
655
656        {
657            let mut recorder = lock_plan(plan);
658            if !recorder.record(&sub_name, StepKind::Workflow, &self.workflow_name, None) {
659                return Ok(planned);
660            }
661            recorder.set_last(vec![sub_name.clone()]);
662        }
663
664        let expand = lock_plan(plan).enter_workflow();
665        if expand {
666            let previous_payload = lock_plan(plan).swap_payload(payload.clone());
667
668            let mut child = WorkflowContext::new(
669                Uuid::now_v7(),
670                sub_name.clone(),
671                self.store.clone(),
672                self.provider.clone(),
673            );
674            child.handler_resolver = self.handler_resolver.clone();
675            child.set_plan(plan.clone());
676
677            if let Err(err) = handler.execute(&mut child).await {
678                lock_plan(plan).fail(format!(
679                    "sub-workflow {sub_name} could not be planned: {err}"
680                ));
681            }
682
683            let mut recorder = lock_plan(plan);
684            recorder.swap_payload(previous_payload);
685            recorder.leave_workflow();
686        }
687
688        Ok(planned)
689    }
690
691    /// Execute a child workflow and return aggregated output plus whether
692    /// at least one `allow_failure` step failed.
693    ///
694    /// When another active run holds the step's concurrency key, no child run
695    /// is created and [`ChildOutcome::Conflict`] is returned.
696    ///
697    /// `resume` is the child run recorded on an open or interrupted step: that
698    /// run is re-entered, with its completed steps replayed, instead of
699    /// creating a new one. A child re-entered after a lost lease has its
700    /// `Running` steps marked interrupted first, like a requeued run, and the
701    /// new step records the same child run so a second interruption re-enters
702    /// it again. A child that suspends is left in its suspension status and
703    /// [`EngineError::ChildSuspended`] is returned.
704    async fn execute_child_workflow(
705        &self,
706        config: &WorkflowStepConfig,
707        step_id: Uuid,
708        resume: Option<ChildResume>,
709    ) -> Result<ChildOutcome, EngineError> {
710        let resolver = self.handler_resolver.as_ref().ok_or_else(|| {
711            EngineError::InvalidWorkflow(
712                "sub-workflow requires a handler resolver (use Engine to execute)".to_string(),
713            )
714        })?;
715
716        let handler = resolver(&config.workflow_name).ok_or_else(|| {
717            EngineError::InvalidWorkflow(format!("no handler registered: {}", config.workflow_name))
718        })?;
719
720        let (child_run_id, carried_cost_usd, carried_duration_ms) = match resume {
721            Some(resume) => {
722                let child_run_id = resume.run_id();
723                if let ChildResume::Interrupted(_) = resume {
724                    self.store
725                        .update_step(
726                            step_id,
727                            StepUpdate {
728                                output: Some(json!({ CHILD_RUN_ID_KEY: child_run_id })),
729                                ..StepUpdate::default()
730                            },
731                        )
732                        .await?;
733                }
734
735                let child_run = self
736                    .store
737                    .get_run(child_run_id)
738                    .await?
739                    .ok_or(EngineError::Store(StoreError::RunNotFound(child_run_id)))?;
740
741                match child_run.status.state {
742                    // Left running by the worker that lost the lease: its open
743                    // steps are executed again, like those of a requeued run.
744                    RunStatus::Running if matches!(resume, ChildResume::Interrupted(_)) => {
745                        interrupt_running_steps(self.store.as_ref(), child_run_id).await?;
746                    }
747                    // Already moved to Running by the path that resumed it.
748                    RunStatus::Running => {}
749                    RunStatus::AwaitingApproval | RunStatus::Pending => {
750                        self.store
751                            .update_run_status(child_run_id, RunStatus::Running)
752                            .await?;
753                    }
754                    RunStatus::Sleeping => {
755                        self.store
756                            .update_run_status(child_run_id, RunStatus::Pending)
757                            .await?;
758                        self.store
759                            .update_run_status(child_run_id, RunStatus::Running)
760                            .await?;
761                    }
762                    // Cancelled while the chain waited, or before the parent
763                    // closed its step: the cancellation is the outcome.
764                    RunStatus::Cancelled => {
765                        return cancelled_child_outcome(
766                            config,
767                            &child_run,
768                            child_run.cost_usd,
769                            child_run.duration_ms,
770                            child_run.output.clone(),
771                        );
772                    }
773                    // The child finished but the parent stopped before closing
774                    // its step: report the recorded outcome, run nothing.
775                    status @ RunStatus::Failed if config.allow_failure => {
776                        let error = child_run
777                            .error
778                            .clone()
779                            .unwrap_or_else(|| "child run failed".to_string());
780                        return Ok(ChildOutcome::Finished(
781                            SubWorkflowOutput::new(
782                                child_run_id,
783                                &config.workflow_name,
784                                status,
785                                child_run.cost_usd,
786                                child_run.duration_ms,
787                            )
788                            .with_output(child_run.output.clone())
789                            .with_error(error),
790                            true,
791                        ));
792                    }
793                    status @ (RunStatus::Completed | RunStatus::Warning) => {
794                        return Ok(ChildOutcome::Finished(
795                            SubWorkflowOutput::new(
796                                child_run_id,
797                                &config.workflow_name,
798                                status,
799                                child_run.cost_usd,
800                                child_run.duration_ms,
801                            )
802                            .with_output(child_run.output.clone()),
803                            status == RunStatus::Warning,
804                        ));
805                    }
806                    other => {
807                        return Err(EngineError::InvalidWorkflow(format!(
808                            "child run {child_run_id} is {other}"
809                        )));
810                    }
811                }
812
813                info!(
814                    parent_run_id = %self.run_id,
815                    child_run_id = %child_run_id,
816                    workflow = %config.workflow_name,
817                    "child run re-entered"
818                );
819                (child_run_id, child_run.cost_usd, child_run.duration_ms)
820            }
821            None => {
822                let priority = config
823                    .priority
824                    .unwrap_or_else(|| clamp_priority(handler.priority()));
825                let child_run_id = match self
826                    .create_child_run(config, priority, handler.as_ref())
827                    .await
828                {
829                    Ok(id) => id,
830                    Err(EngineError::ConcurrencyConflict { key, run_id }) => {
831                        return Ok(ChildOutcome::Conflict(ConcurrencyConflict::new(
832                            key, run_id,
833                        )));
834                    }
835                    Err(err) => return Err(err),
836                };
837
838                // Recorded before the child runs, so a suspension of the child
839                // can be resumed into this same run.
840                self.store
841                    .update_step(
842                        step_id,
843                        StepUpdate {
844                            output: Some(json!({ CHILD_RUN_ID_KEY: child_run_id })),
845                            ..StepUpdate::default()
846                        },
847                    )
848                    .await?;
849
850                self.store
851                    .update_run_status(child_run_id, RunStatus::Running)
852                    .await?;
853                (child_run_id, Decimal::ZERO, 0)
854            }
855        };
856
857        let run_start = Instant::now();
858        let mut child_ctx = WorkflowContext {
859            run_id: child_run_id,
860            root_run_id: self.root_run_id,
861            workflow_name: config.workflow_name.clone(),
862            store: self.store.clone(),
863            provider: self.provider.clone(),
864            decision_provider: self.decision_provider.clone(),
865            handler_resolver: self.handler_resolver.clone(),
866            position: 0,
867            last_step_ids: Vec::new(),
868            // A re-entered child starts from what it already spent, like a
869            // resumed top-level run.
870            total_cost_usd: carried_cost_usd,
871            total_duration_ms: 0,
872            max_cost_usd: self.max_cost_usd,
873            // Everything the parent chain already spent counts against the
874            // shared cap, so the child cannot restart the budget from zero.
875            inherited_cost_usd: self.charged_cost_usd(),
876            replay_steps: HashMap::new(),
877            replay_wave_steps: HashMap::new(),
878            granted_approvals: HashMap::new(),
879            answered_inputs: HashMap::new(),
880            interrupted_children: HashMap::new(),
881            // A child run is never itself retried.
882            attempt: 1,
883            carried_duration_ms,
884            log_sender: self.log_sender.clone(),
885            // A child shares the storage backend but not the parent's artifacts:
886            // input lookups are scoped to the child's own run.
887            artifact_sink: self.artifact_sink.clone(),
888            has_allowed_failure: false,
889            error_handlers: Vec::new(),
890            guard_state: self.guard_state.clone(),
891            guard_config: self.guard_config.clone(),
892            step_results: Vec::new(),
893            event_bus: self.event_bus.clone(),
894            // A child run is mocked exactly like its parent.
895            interceptor: self.interceptor.clone(),
896            trace_context: self.trace_context.child(),
897            operation_ctx: None,
898            run_created_at: None,
899            plan: None,
900            output: None,
901            worker_tags: self.worker_tags.clone(),
902        };
903
904        // A re-entered child replays its completed steps and is served the
905        // answer, signal or elapsed delay it was suspended on.
906        let loaded = if resume.is_some() {
907            child_ctx.load_replay_steps().await
908        } else {
909            Ok(())
910        };
911        let result = match loaded {
912            Ok(()) => handler.execute(&mut child_ctx).await,
913            Err(err) => Err(err),
914        };
915        let total_duration = child_ctx.carried_duration_ms + run_start.elapsed().as_millis() as u64;
916        let completed_at = Utc::now();
917
918        // Cancelled while it ran: whatever the handler returned, the child
919        // stays cancelled and the cancellation is its outcome.
920        if let Some(child_run) = self.store.get_run(child_run_id).await? {
921            match child_run.status.state {
922                RunStatus::Cancelled => {
923                    return cancelled_child_outcome(
924                        config,
925                        &child_run,
926                        child_ctx.total_cost_usd,
927                        total_duration,
928                        child_ctx.output().cloned(),
929                    );
930                }
931                // Paused with its root while it ran: the child is left as the
932                // pause recorded it and the parent stops with it.
933                RunStatus::Paused => {
934                    info!(
935                        parent_run_id = %self.run_id,
936                        child_run_id = %child_run_id,
937                        "child run paused"
938                    );
939                    return Err(EngineError::RunPaused {
940                        run_id: child_run_id,
941                    });
942                }
943                _ => {}
944            }
945        }
946
947        match result {
948            Ok(()) => {
949                let child_status = if child_ctx.has_allowed_failure {
950                    RunStatus::Warning
951                } else {
952                    RunStatus::Completed
953                };
954                self.store
955                    .update_run(
956                        child_run_id,
957                        RunUpdate {
958                            status: Some(child_status),
959                            cost_usd: Some(child_ctx.total_cost_usd),
960                            duration_ms: Some(total_duration),
961                            completed_at: Some(completed_at),
962                            output: child_ctx.output().cloned(),
963                            ..RunUpdate::default()
964                        },
965                    )
966                    .await?;
967
968                let child_had_allowed_failure = child_ctx.has_allowed_failure;
969                Ok(ChildOutcome::Finished(
970                    SubWorkflowOutput::new(
971                        child_run_id,
972                        &config.workflow_name,
973                        child_status,
974                        child_ctx.total_cost_usd,
975                        total_duration,
976                    )
977                    .with_output(child_ctx.output().cloned()),
978                    child_had_allowed_failure,
979                ))
980            }
981            Err(err) if err.is_suspension() => {
982                match self
983                    .suspend_child_run(child_run_id, &err, child_ctx.total_cost_usd, total_duration)
984                    .await
985                {
986                    Ok(()) => {
987                        info!(
988                            parent_run_id = %self.run_id,
989                            child_run_id = %child_run_id,
990                            cause = %err.suspension_leaf(),
991                            "child run suspended"
992                        );
993                        Err(EngineError::ChildSuspended {
994                            run_id: child_run_id,
995                            cause: Box::new(err),
996                        })
997                    }
998                    Err(store_err) => {
999                        self.fail_child_run(
1000                            child_run_id,
1001                            RunStatus::Failed,
1002                            &store_err,
1003                            child_ctx.total_cost_usd,
1004                            total_duration,
1005                            child_ctx.output().cloned(),
1006                        )
1007                        .await;
1008                        Err(store_err)
1009                    }
1010                }
1011            }
1012            Err(err) => {
1013                // The engine cancels top-level runs stopped by a guardrail.
1014                let status = if matches!(
1015                    err,
1016                    EngineError::RunBudgetExceeded { .. } | EngineError::WorkflowGuardRejected(_)
1017                ) {
1018                    RunStatus::Cancelled
1019                } else {
1020                    RunStatus::Failed
1021                };
1022                self.fail_child_run(
1023                    child_run_id,
1024                    status,
1025                    &err,
1026                    child_ctx.total_cost_usd,
1027                    total_duration,
1028                    child_ctx.output().cloned(),
1029                )
1030                .await;
1031                if config.allow_failure {
1032                    return Ok(ChildOutcome::Finished(
1033                        SubWorkflowOutput::new(
1034                            child_run_id,
1035                            &config.workflow_name,
1036                            status,
1037                            child_ctx.total_cost_usd,
1038                            total_duration,
1039                        )
1040                        .with_output(child_ctx.output().cloned())
1041                        .with_error(err.to_string()),
1042                        true,
1043                    ));
1044                }
1045                Err(err)
1046            }
1047        }
1048    }
1049
1050    /// Create the child run of a sub-workflow step and return its id.
1051    ///
1052    /// The child inherits the parent labels and author, and is linked to its
1053    /// parent and to the root of the chain by two labels, so a suspended child
1054    /// can be found and resumed like a top-level run.
1055    async fn create_child_run(
1056        &self,
1057        config: &WorkflowStepConfig,
1058        priority: i16,
1059        handler: &dyn WorkflowHandler,
1060    ) -> Result<Uuid, EngineError> {
1061        // Whoever triggered the parent workflow is accountable for its children.
1062        let parent = self.store.get_run(self.run_id).await?;
1063        let (mut labels, parent_author) =
1064            parent.map(|r| (r.labels, r.created_by)).unwrap_or_default();
1065        // Overwritten, never inherited: a grand-child must point at its own
1066        // parent, not at its grand-parent.
1067        labels.insert(PARENT_RUN_ID_LABEL.to_string(), self.run_id.to_string());
1068        labels.insert(LABEL_ROOT_RUN_ID.to_string(), self.root_run_id.to_string());
1069
1070        let child_run = self
1071            .store
1072            .create_run(NewRun {
1073                workflow_name: config.workflow_name.clone(),
1074                trigger: TriggerKind::Workflow,
1075                payload: config.payload.clone(),
1076                max_retries: 0,
1077                handler_version: None,
1078                labels,
1079                scheduled_at: None,
1080                created_by: parent_author,
1081                idempotency_key: None,
1082                concurrency_key: config.concurrency_key.clone(),
1083                priority,
1084                // A child runs inside its parent's slot: it never consumes a
1085                // concurrency group slot of its own.
1086                concurrency_limits: Vec::new(),
1087                // The child shares the parent's cap; it does not get its own budget.
1088                max_cost_usd: self.max_cost_usd,
1089                // Recorded for display: the child runs on its parent's worker.
1090                worker_tags: normalize_worker_tags(handler.required_worker_tags()),
1091            })
1092            .await?
1093            .into_run();
1094
1095        info!(
1096            parent_run_id = %self.run_id,
1097            child_run_id = %child_run.id,
1098            workflow = %config.workflow_name,
1099            "child run created"
1100        );
1101        Ok(child_run.id)
1102    }
1103
1104    /// Persist the suspension of a child run, with no event: the root run
1105    /// publishes the suspension once the whole chain is suspended.
1106    ///
1107    /// A direct suspension is persisted like a top-level run's (a delay, a
1108    /// capacity wait or a signal deadline arms `scheduled_at`; a capacity wait
1109    /// also records its provider kind). A child suspended because of its own
1110    /// child gets no `scheduled_at`: only the deepest run owns the
1111    /// wake-up, so the chain is never resumed twice.
1112    async fn suspend_child_run(
1113        &self,
1114        child_run_id: Uuid,
1115        err: &EngineError,
1116        cost_usd: Decimal,
1117        duration_ms: u64,
1118    ) -> Result<(), EngineError> {
1119        let totals = RunUpdate {
1120            cost_usd: Some(cost_usd),
1121            duration_ms: Some(duration_ms),
1122            ..RunUpdate::default()
1123        };
1124
1125        let update = match err {
1126            EngineError::DelaySleeping { wake_at, .. } => RunUpdate {
1127                status: Some(RunStatus::Sleeping),
1128                scheduled_at: Some(*wake_at),
1129                ..totals
1130            },
1131            EngineError::CapacitySleeping { kind, wake_at, .. } => RunUpdate {
1132                status: Some(RunStatus::Sleeping),
1133                scheduled_at: Some(*wake_at),
1134                capacity_wait_kind: Some(ProviderKind::new(kind.as_str())),
1135                ..totals
1136            },
1137            EngineError::SignalWaiting {
1138                step_id,
1139                deadline_at,
1140                ..
1141            } => {
1142                // Atomic with the step lock, like a top-level run.
1143                self.store
1144                    .suspend_run_on_signal(child_run_id, *step_id, *deadline_at)
1145                    .await?;
1146                totals
1147            }
1148            EngineError::ChildSuspended { cause, .. } => RunUpdate {
1149                status: Some(cause.suspension_status()),
1150                ..totals
1151            },
1152            _ => RunUpdate {
1153                status: Some(RunStatus::AwaitingApproval),
1154                ..totals
1155            },
1156        };
1157
1158        self.store.update_run(child_run_id, update).await?;
1159        Ok(())
1160    }
1161
1162    /// Mark a child run failed after its handler (or its suspension) failed.
1163    ///
1164    /// Best effort: the original error is what the parent reports.
1165    async fn fail_child_run(
1166        &self,
1167        child_run_id: Uuid,
1168        status: RunStatus,
1169        err: &EngineError,
1170        cost_usd: Decimal,
1171        duration_ms: u64,
1172        output: Option<Value>,
1173    ) {
1174        if let Err(store_err) = self
1175            .store
1176            .update_run(
1177                child_run_id,
1178                RunUpdate {
1179                    status: Some(status),
1180                    error: Some(err.to_string()),
1181                    cost_usd: Some(cost_usd),
1182                    duration_ms: Some(duration_ms),
1183                    completed_at: Some(Utc::now()),
1184                    output,
1185                    ..RunUpdate::default()
1186                },
1187            )
1188            .await
1189        {
1190            error!(
1191                child_run_id = %child_run_id,
1192                store_error = %store_err,
1193                "failed to persist child run failure"
1194            );
1195        }
1196    }
1197}