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