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