Skip to main content

ironflow_engine/context/steps/
sub_workflow.rs

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