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