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};
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 config.concurrency_key = options.into_concurrency_key();
379 let position = self.position;
380
381 let existing = self.replay_steps.get(&position).cloned();
382 if let Some(existing) = &existing {
383 check_replay_identity(
384 existing,
385 position,
386 &config.workflow_name,
387 &StepKind::Workflow,
388 )?;
389 if existing.status.state == StepStatus::Completed {
390 return self.replay_sub_workflow(existing);
391 }
392 }
393
394 // A step interrupted by a lost lease must be the same step the handler
395 // calls now before its child run is re-entered.
396 let interrupted = self.interrupted_children.get(&position).cloned();
397 if let Some(interrupted) = &interrupted {
398 check_replay_identity(
399 interrupted,
400 position,
401 &config.workflow_name,
402 &StepKind::Workflow,
403 )?;
404 }
405
406 // A child runs on the worker that runs its parent: refuse it here
407 // rather than run it on a host missing what it requires.
408 self.check_child_worker_tags(handler)?;
409
410 // Guard check: verify limits before creating the step.
411 if let (Some(guard_config), Some(guard_state)) = (&self.guard_config, &self.guard_state) {
412 let state = guard_state
413 .lock()
414 .map_err(|_| WorkflowRejection::GuardUnavailable)?;
415 state.check(guard_config, handler.name())?;
416 }
417
418 self.position += 1;
419
420 // An open step was left by a child that suspended: reuse it instead of
421 // recording a second step at the same position.
422 let (step, resume) = match existing.filter(|s| s.status.state == StepStatus::Running) {
423 Some(step) => {
424 let resume = recorded_child_run_id(&step).map(ChildResume::Suspended);
425 (step, resume)
426 }
427 None => {
428 let trace_id = step_trace_id(self.run_id, &config.workflow_name, position);
429 let step = self
430 .store
431 .create_step(NewStep {
432 run_id: self.run_id,
433 trace_id,
434 name: config.workflow_name.clone(),
435 kind: StepKind::Workflow,
436 position,
437 input: Some(to_value(&config)?),
438 is_error_handler: false,
439 })
440 .await?;
441
442 self.start_step(step.id, Utc::now()).await?;
443 // A step interrupted by a lost lease re-enters the child run it
444 // recorded instead of starting a new one.
445 let resume = interrupted
446 .as_ref()
447 .and_then(recorded_child_run_id)
448 .map(ChildResume::Interrupted);
449 (step, resume)
450 }
451 };
452
453 // Record invocation in guard state (fail-closed).
454 if let Some(guard_state) = &self.guard_state {
455 let mut state = guard_state
456 .lock()
457 .map_err(|_| WorkflowRejection::GuardUnavailable)?;
458 state.record_invocation(handler.name());
459 }
460
461 match self.execute_child_workflow(&config, step.id, resume).await {
462 // No child run was created: the step completes with the conflict
463 // as its output, so a replay serves the same outcome.
464 Ok(ChildOutcome::Conflict(conflict)) => {
465 self.store
466 .update_step(
467 step.id,
468 StepUpdate {
469 status: Some(StepStatus::Completed),
470 output: Some(json!({ "concurrency_conflict": conflict })),
471 duration_ms: Some(0),
472 cost_usd: Some(Decimal::ZERO),
473 completed_at: Some(Utc::now()),
474 ..StepUpdate::default()
475 },
476 )
477 .await?;
478
479 info!(
480 run_id = %self.run_id,
481 child_workflow = %config.workflow_name,
482 key = %conflict.key(),
483 holder = %conflict.run_id(),
484 "workflow step skipped: concurrency conflict"
485 );
486
487 self.last_step_ids = vec![step.id];
488
489 self.guard_record_return();
490 Ok(SubWorkflowOutcome::Conflict(conflict))
491 }
492 Ok(ChildOutcome::Finished(output, child_had_allowed_failure)) => {
493 self.total_cost_usd += output.cost_usd();
494 self.total_duration_ms += output.duration_ms();
495 if child_had_allowed_failure {
496 self.has_allowed_failure = true;
497 }
498
499 let completed_at = Utc::now();
500 self.store
501 .update_step(
502 step.id,
503 StepUpdate {
504 status: Some(StepStatus::Completed),
505 output: Some(to_value(&output)?),
506 duration_ms: Some(output.duration_ms()),
507 cost_usd: Some(output.cost_usd()),
508 completed_at: Some(completed_at),
509 ..StepUpdate::default()
510 },
511 )
512 .await?;
513
514 info!(
515 run_id = %self.run_id,
516 child_workflow = %config.workflow_name,
517 duration_ms = output.duration_ms(),
518 "workflow step completed"
519 );
520
521 self.last_step_ids = vec![step.id];
522
523 self.guard_record_return();
524 Ok(SubWorkflowOutcome::Completed(output))
525 }
526 // The child suspended: the step stays open, neither failed nor
527 // completed, so the next replay re-enters the same child run.
528 Err(err) if err.is_suspension() => {
529 self.guard_record_return();
530 Err(err)
531 }
532 Err(err) => {
533 let completed_at = Utc::now();
534 if let Err(store_err) = self
535 .store
536 .update_step(
537 step.id,
538 StepUpdate {
539 status: Some(StepStatus::Failed),
540 error: Some(err.to_string()),
541 completed_at: Some(completed_at),
542 ..StepUpdate::default()
543 },
544 )
545 .await
546 {
547 error!(step_id = %step.id, error = %store_err, "failed to persist step failure");
548 }
549
550 self.guard_record_return();
551 Err(err)
552 }
553 }
554 }
555
556 /// Replay a `Workflow` step completed in a previous execution: the child
557 /// run is not executed again and nothing is re-counted by the guard.
558 ///
559 /// A step skipped on a concurrency conflict replays the same conflict: the
560 /// child is not attempted again, even if the key has been released since.
561 fn replay_sub_workflow(&mut self, step: &Step) -> Result<SubWorkflowOutcome, EngineError> {
562 let recorded = step.output.clone().ok_or_else(|| {
563 EngineError::StepConfig(format!(
564 "completed workflow step {} has no recorded output",
565 step.id
566 ))
567 })?;
568 let recorded: RecordedWorkflowStep = from_value(recorded)?;
569
570 self.position += 1;
571 self.last_step_ids = vec![step.id];
572
573 let output = match SubWorkflowOutcome::from(recorded) {
574 SubWorkflowOutcome::Completed(output) => output,
575 SubWorkflowOutcome::Conflict(conflict) => {
576 info!(
577 run_id = %self.run_id,
578 step = %step.name,
579 key = %conflict.key(),
580 holder = %conflict.run_id(),
581 "workflow step replayed: concurrency conflict"
582 );
583 return Ok(SubWorkflowOutcome::Conflict(conflict));
584 }
585 };
586
587 // Cost is not added: `carry_over_run_totals` seeded `total_cost_usd`
588 // from the run totals persisted before the suspension, which already
589 // include this child.
590 self.total_duration_ms += output.duration_ms();
591 if matches!(
592 output.status(),
593 RunStatus::Warning | RunStatus::Failed | RunStatus::Cancelled
594 ) {
595 self.has_allowed_failure = true;
596 }
597
598 info!(
599 run_id = %self.run_id,
600 child_run_id = %output.run_id(),
601 step = %step.name,
602 "workflow step replayed from previous execution"
603 );
604 Ok(SubWorkflowOutcome::Completed(output))
605 }
606
607 /// Refuse a child whose required worker tags are not all carried by the
608 /// worker running this context.
609 ///
610 /// No check is made outside a tagged worker (API or local mode).
611 fn check_child_worker_tags(&self, handler: &dyn WorkflowHandler) -> Result<(), EngineError> {
612 let Some(carried) = &self.worker_tags else {
613 return Ok(());
614 };
615 let missing: Vec<String> = normalize_worker_tags(handler.required_worker_tags())
616 .into_iter()
617 .filter(|tag| !carried.contains(tag))
618 .collect();
619 if missing.is_empty() {
620 return Ok(());
621 }
622 Err(EngineError::InvalidWorkflow(format!(
623 "sub-workflow '{}' requires worker tags [{}] that this worker does not carry",
624 handler.name(),
625 missing.join(", ")
626 )))
627 }
628
629 /// Record a sub-workflow invocation while planning, expanding the child
630 /// handler into the same plan when the depth limit allows it.
631 ///
632 /// The child plans against its own payload and under its own workflow
633 /// name; the parent's payload is restored on the way out.
634 async fn plan_sub_workflow(
635 &mut self,
636 plan: &SharedPlanRecorder,
637 handler: &dyn WorkflowHandler,
638 payload: Value,
639 ) -> Result<SubWorkflowOutput, EngineError> {
640 self.position += 1;
641 let sub_name = handler.name().to_string();
642 // No child run exists while planning: a nil id and zero metrics.
643 let planned = SubWorkflowOutput::new(
644 Uuid::nil(),
645 &sub_name,
646 RunStatus::Completed,
647 Decimal::ZERO,
648 0,
649 );
650
651 {
652 let mut recorder = lock_plan(plan);
653 if !recorder.record(&sub_name, StepKind::Workflow, &self.workflow_name, None) {
654 return Ok(planned);
655 }
656 recorder.set_last(vec![sub_name.clone()]);
657 }
658
659 let expand = lock_plan(plan).enter_workflow();
660 if expand {
661 let previous_payload = lock_plan(plan).swap_payload(payload.clone());
662
663 let mut child = WorkflowContext::new(
664 Uuid::now_v7(),
665 sub_name.clone(),
666 self.store.clone(),
667 self.provider.clone(),
668 );
669 child.handler_resolver = self.handler_resolver.clone();
670 child.set_plan(plan.clone());
671
672 if let Err(err) = handler.execute(&mut child).await {
673 lock_plan(plan).fail(format!(
674 "sub-workflow {sub_name} could not be planned: {err}"
675 ));
676 }
677
678 let mut recorder = lock_plan(plan);
679 recorder.swap_payload(previous_payload);
680 recorder.leave_workflow();
681 }
682
683 Ok(planned)
684 }
685
686 /// Execute a child workflow and return aggregated output plus whether
687 /// at least one `allow_failure` step failed.
688 ///
689 /// When another active run holds the step's concurrency key, no child run
690 /// is created and [`ChildOutcome::Conflict`] is returned.
691 ///
692 /// `resume` is the child run recorded on an open or interrupted step: that
693 /// run is re-entered, with its completed steps replayed, instead of
694 /// creating a new one. A child re-entered after a lost lease has its
695 /// `Running` steps marked interrupted first, like a requeued run, and the
696 /// new step records the same child run so a second interruption re-enters
697 /// it again. A child that suspends is left in its suspension status and
698 /// [`EngineError::ChildSuspended`] is returned.
699 async fn execute_child_workflow(
700 &self,
701 config: &WorkflowStepConfig,
702 step_id: Uuid,
703 resume: Option<ChildResume>,
704 ) -> Result<ChildOutcome, EngineError> {
705 let resolver = self.handler_resolver.as_ref().ok_or_else(|| {
706 EngineError::InvalidWorkflow(
707 "sub-workflow requires a handler resolver (use Engine to execute)".to_string(),
708 )
709 })?;
710
711 let handler = resolver(&config.workflow_name).ok_or_else(|| {
712 EngineError::InvalidWorkflow(format!("no handler registered: {}", config.workflow_name))
713 })?;
714
715 let (child_run_id, carried_cost_usd, carried_duration_ms) = match resume {
716 Some(resume) => {
717 let child_run_id = resume.run_id();
718 if let ChildResume::Interrupted(_) = resume {
719 self.store
720 .update_step(
721 step_id,
722 StepUpdate {
723 output: Some(json!({ CHILD_RUN_ID_KEY: child_run_id })),
724 ..StepUpdate::default()
725 },
726 )
727 .await?;
728 }
729
730 let child_run = self
731 .store
732 .get_run(child_run_id)
733 .await?
734 .ok_or(EngineError::Store(StoreError::RunNotFound(child_run_id)))?;
735
736 match child_run.status.state {
737 // Left running by the worker that lost the lease: its open
738 // steps are executed again, like those of a requeued run.
739 RunStatus::Running if matches!(resume, ChildResume::Interrupted(_)) => {
740 interrupt_running_steps(self.store.as_ref(), child_run_id).await?;
741 }
742 // Already moved to Running by the path that resumed it.
743 RunStatus::Running => {}
744 RunStatus::AwaitingApproval | RunStatus::Pending => {
745 self.store
746 .update_run_status(child_run_id, RunStatus::Running)
747 .await?;
748 }
749 RunStatus::Sleeping => {
750 self.store
751 .update_run_status(child_run_id, RunStatus::Pending)
752 .await?;
753 self.store
754 .update_run_status(child_run_id, RunStatus::Running)
755 .await?;
756 }
757 // Cancelled while the chain waited, or before the parent
758 // closed its step: the cancellation is the outcome.
759 RunStatus::Cancelled => {
760 return cancelled_child_outcome(
761 config,
762 &child_run,
763 child_run.cost_usd,
764 child_run.duration_ms,
765 child_run.output.clone(),
766 );
767 }
768 // The child finished but the parent stopped before closing
769 // its step: report the recorded outcome, run nothing.
770 status @ RunStatus::Failed if config.allow_failure => {
771 let error = child_run
772 .error
773 .clone()
774 .unwrap_or_else(|| "child run failed".to_string());
775 return Ok(ChildOutcome::Finished(
776 SubWorkflowOutput::new(
777 child_run_id,
778 &config.workflow_name,
779 status,
780 child_run.cost_usd,
781 child_run.duration_ms,
782 )
783 .with_output(child_run.output.clone())
784 .with_error(error),
785 true,
786 ));
787 }
788 status @ (RunStatus::Completed | RunStatus::Warning) => {
789 return Ok(ChildOutcome::Finished(
790 SubWorkflowOutput::new(
791 child_run_id,
792 &config.workflow_name,
793 status,
794 child_run.cost_usd,
795 child_run.duration_ms,
796 )
797 .with_output(child_run.output.clone()),
798 status == RunStatus::Warning,
799 ));
800 }
801 other => {
802 return Err(EngineError::InvalidWorkflow(format!(
803 "child run {child_run_id} is {other}"
804 )));
805 }
806 }
807
808 info!(
809 parent_run_id = %self.run_id,
810 child_run_id = %child_run_id,
811 workflow = %config.workflow_name,
812 "child run re-entered"
813 );
814 (child_run_id, child_run.cost_usd, child_run.duration_ms)
815 }
816 None => {
817 let child_run_id = match self.create_child_run(config, handler.as_ref()).await {
818 Ok(id) => id,
819 Err(EngineError::ConcurrencyConflict { key, run_id }) => {
820 return Ok(ChildOutcome::Conflict(ConcurrencyConflict::new(
821 key, run_id,
822 )));
823 }
824 Err(err) => return Err(err),
825 };
826
827 // Recorded before the child runs, so a suspension of the child
828 // can be resumed into this same run.
829 self.store
830 .update_step(
831 step_id,
832 StepUpdate {
833 output: Some(json!({ CHILD_RUN_ID_KEY: child_run_id })),
834 ..StepUpdate::default()
835 },
836 )
837 .await?;
838
839 self.store
840 .update_run_status(child_run_id, RunStatus::Running)
841 .await?;
842 (child_run_id, Decimal::ZERO, 0)
843 }
844 };
845
846 let run_start = Instant::now();
847 let mut child_ctx = WorkflowContext {
848 run_id: child_run_id,
849 root_run_id: self.root_run_id,
850 workflow_name: config.workflow_name.clone(),
851 store: self.store.clone(),
852 provider: self.provider.clone(),
853 decision_provider: self.decision_provider.clone(),
854 handler_resolver: self.handler_resolver.clone(),
855 position: 0,
856 last_step_ids: Vec::new(),
857 // A re-entered child starts from what it already spent, like a
858 // resumed top-level run.
859 total_cost_usd: carried_cost_usd,
860 total_duration_ms: 0,
861 max_cost_usd: self.max_cost_usd,
862 // Everything the parent chain already spent counts against the
863 // shared cap, so the child cannot restart the budget from zero.
864 inherited_cost_usd: self.charged_cost_usd(),
865 replay_steps: HashMap::new(),
866 replay_wave_steps: HashMap::new(),
867 granted_approvals: HashMap::new(),
868 answered_inputs: HashMap::new(),
869 interrupted_children: HashMap::new(),
870 // A child run is never itself retried.
871 attempt: 1,
872 carried_duration_ms,
873 log_sender: self.log_sender.clone(),
874 // A child shares the storage backend but not the parent's artifacts:
875 // input lookups are scoped to the child's own run.
876 artifact_sink: self.artifact_sink.clone(),
877 has_allowed_failure: false,
878 error_handlers: Vec::new(),
879 guard_state: self.guard_state.clone(),
880 guard_config: self.guard_config.clone(),
881 step_results: Vec::new(),
882 event_bus: self.event_bus.clone(),
883 // A child run is mocked exactly like its parent.
884 interceptor: self.interceptor.clone(),
885 trace_context: self.trace_context.child(),
886 operation_ctx: None,
887 run_created_at: None,
888 plan: None,
889 output: None,
890 worker_tags: self.worker_tags.clone(),
891 };
892
893 // A re-entered child replays its completed steps and is served the
894 // answer, signal or elapsed delay it was suspended on.
895 let loaded = if resume.is_some() {
896 child_ctx.load_replay_steps().await
897 } else {
898 Ok(())
899 };
900 let result = match loaded {
901 Ok(()) => handler.execute(&mut child_ctx).await,
902 Err(err) => Err(err),
903 };
904 let total_duration = child_ctx.carried_duration_ms + run_start.elapsed().as_millis() as u64;
905 let completed_at = Utc::now();
906
907 // Cancelled while it ran: whatever the handler returned, the child
908 // stays cancelled and the cancellation is its outcome.
909 if let Some(child_run) = self.store.get_run(child_run_id).await?
910 && child_run.status.state == RunStatus::Cancelled
911 {
912 return cancelled_child_outcome(
913 config,
914 &child_run,
915 child_ctx.total_cost_usd,
916 total_duration,
917 child_ctx.output().cloned(),
918 );
919 }
920
921 match result {
922 Ok(()) => {
923 let child_status = if child_ctx.has_allowed_failure {
924 RunStatus::Warning
925 } else {
926 RunStatus::Completed
927 };
928 self.store
929 .update_run(
930 child_run_id,
931 RunUpdate {
932 status: Some(child_status),
933 cost_usd: Some(child_ctx.total_cost_usd),
934 duration_ms: Some(total_duration),
935 completed_at: Some(completed_at),
936 output: child_ctx.output().cloned(),
937 ..RunUpdate::default()
938 },
939 )
940 .await?;
941
942 let child_had_allowed_failure = child_ctx.has_allowed_failure;
943 Ok(ChildOutcome::Finished(
944 SubWorkflowOutput::new(
945 child_run_id,
946 &config.workflow_name,
947 child_status,
948 child_ctx.total_cost_usd,
949 total_duration,
950 )
951 .with_output(child_ctx.output().cloned()),
952 child_had_allowed_failure,
953 ))
954 }
955 Err(err) if err.is_suspension() => {
956 match self
957 .suspend_child_run(child_run_id, &err, child_ctx.total_cost_usd, total_duration)
958 .await
959 {
960 Ok(()) => {
961 info!(
962 parent_run_id = %self.run_id,
963 child_run_id = %child_run_id,
964 cause = %err.suspension_leaf(),
965 "child run suspended"
966 );
967 Err(EngineError::ChildSuspended {
968 run_id: child_run_id,
969 cause: Box::new(err),
970 })
971 }
972 Err(store_err) => {
973 self.fail_child_run(
974 child_run_id,
975 RunStatus::Failed,
976 &store_err,
977 child_ctx.total_cost_usd,
978 total_duration,
979 child_ctx.output().cloned(),
980 )
981 .await;
982 Err(store_err)
983 }
984 }
985 }
986 Err(err) => {
987 // The engine cancels top-level runs stopped by a guardrail.
988 let status = if matches!(
989 err,
990 EngineError::RunBudgetExceeded { .. } | EngineError::WorkflowGuardRejected(_)
991 ) {
992 RunStatus::Cancelled
993 } else {
994 RunStatus::Failed
995 };
996 self.fail_child_run(
997 child_run_id,
998 status,
999 &err,
1000 child_ctx.total_cost_usd,
1001 total_duration,
1002 child_ctx.output().cloned(),
1003 )
1004 .await;
1005 if config.allow_failure {
1006 return Ok(ChildOutcome::Finished(
1007 SubWorkflowOutput::new(
1008 child_run_id,
1009 &config.workflow_name,
1010 status,
1011 child_ctx.total_cost_usd,
1012 total_duration,
1013 )
1014 .with_output(child_ctx.output().cloned())
1015 .with_error(err.to_string()),
1016 true,
1017 ));
1018 }
1019 Err(err)
1020 }
1021 }
1022 }
1023
1024 /// Create the child run of a sub-workflow step and return its id.
1025 ///
1026 /// The child inherits the parent labels and author, and is linked to its
1027 /// parent and to the root of the chain by two labels, so a suspended child
1028 /// can be found and resumed like a top-level run.
1029 async fn create_child_run(
1030 &self,
1031 config: &WorkflowStepConfig,
1032 handler: &dyn WorkflowHandler,
1033 ) -> Result<Uuid, EngineError> {
1034 // Whoever triggered the parent workflow is accountable for its children.
1035 let parent = self.store.get_run(self.run_id).await?;
1036 let (mut labels, parent_author) =
1037 parent.map(|r| (r.labels, r.created_by)).unwrap_or_default();
1038 // Overwritten, never inherited: a grand-child must point at its own
1039 // parent, not at its grand-parent.
1040 labels.insert(PARENT_RUN_ID_LABEL.to_string(), self.run_id.to_string());
1041 labels.insert(LABEL_ROOT_RUN_ID.to_string(), self.root_run_id.to_string());
1042
1043 let child_run = self
1044 .store
1045 .create_run(NewRun {
1046 workflow_name: config.workflow_name.clone(),
1047 trigger: TriggerKind::Workflow,
1048 payload: config.payload.clone(),
1049 max_retries: 0,
1050 handler_version: None,
1051 labels,
1052 scheduled_at: None,
1053 created_by: parent_author,
1054 idempotency_key: None,
1055 concurrency_key: config.concurrency_key.clone(),
1056 // A child runs inside its parent's slot: it never consumes a
1057 // concurrency group slot of its own.
1058 concurrency_limits: Vec::new(),
1059 // The child shares the parent's cap; it does not get its own budget.
1060 max_cost_usd: self.max_cost_usd,
1061 // Recorded for display: the child runs on its parent's worker.
1062 worker_tags: normalize_worker_tags(handler.required_worker_tags()),
1063 })
1064 .await?
1065 .into_run();
1066
1067 info!(
1068 parent_run_id = %self.run_id,
1069 child_run_id = %child_run.id,
1070 workflow = %config.workflow_name,
1071 "child run created"
1072 );
1073 Ok(child_run.id)
1074 }
1075
1076 /// Persist the suspension of a child run, with no event: the root run
1077 /// publishes the suspension once the whole chain is suspended.
1078 ///
1079 /// A direct suspension is persisted like a top-level run's (a delay, a
1080 /// capacity wait or a signal deadline arms `scheduled_at`; a capacity wait
1081 /// also records its provider kind). A child suspended because of its own
1082 /// child gets no `scheduled_at`: only the deepest run owns the
1083 /// wake-up, so the chain is never resumed twice.
1084 async fn suspend_child_run(
1085 &self,
1086 child_run_id: Uuid,
1087 err: &EngineError,
1088 cost_usd: Decimal,
1089 duration_ms: u64,
1090 ) -> Result<(), EngineError> {
1091 let totals = RunUpdate {
1092 cost_usd: Some(cost_usd),
1093 duration_ms: Some(duration_ms),
1094 ..RunUpdate::default()
1095 };
1096
1097 let update = match err {
1098 EngineError::DelaySleeping { wake_at, .. } => RunUpdate {
1099 status: Some(RunStatus::Sleeping),
1100 scheduled_at: Some(*wake_at),
1101 ..totals
1102 },
1103 EngineError::CapacitySleeping { kind, wake_at, .. } => RunUpdate {
1104 status: Some(RunStatus::Sleeping),
1105 scheduled_at: Some(*wake_at),
1106 capacity_wait_kind: Some(ProviderKind::new(kind.as_str())),
1107 ..totals
1108 },
1109 EngineError::SignalWaiting {
1110 step_id,
1111 deadline_at,
1112 ..
1113 } => {
1114 // Atomic with the step lock, like a top-level run.
1115 self.store
1116 .suspend_run_on_signal(child_run_id, *step_id, *deadline_at)
1117 .await?;
1118 totals
1119 }
1120 EngineError::ChildSuspended { cause, .. } => RunUpdate {
1121 status: Some(cause.suspension_status()),
1122 ..totals
1123 },
1124 _ => RunUpdate {
1125 status: Some(RunStatus::AwaitingApproval),
1126 ..totals
1127 },
1128 };
1129
1130 self.store.update_run(child_run_id, update).await?;
1131 Ok(())
1132 }
1133
1134 /// Mark a child run failed after its handler (or its suspension) failed.
1135 ///
1136 /// Best effort: the original error is what the parent reports.
1137 async fn fail_child_run(
1138 &self,
1139 child_run_id: Uuid,
1140 status: RunStatus,
1141 err: &EngineError,
1142 cost_usd: Decimal,
1143 duration_ms: u64,
1144 output: Option<Value>,
1145 ) {
1146 if let Err(store_err) = self
1147 .store
1148 .update_run(
1149 child_run_id,
1150 RunUpdate {
1151 status: Some(status),
1152 error: Some(err.to_string()),
1153 cost_usd: Some(cost_usd),
1154 duration_ms: Some(duration_ms),
1155 completed_at: Some(Utc::now()),
1156 output,
1157 ..RunUpdate::default()
1158 },
1159 )
1160 .await
1161 {
1162 error!(
1163 child_run_id = %child_run_id,
1164 store_error = %store_err,
1165 "failed to persist child run failure"
1166 );
1167 }
1168 }
1169}