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}