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