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