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