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