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
7use std::collections::HashMap;
8use std::time::Instant;
9
10use chrono::Utc;
11use rust_decimal::Decimal;
12use serde_json::{Value, to_value};
13use tracing::{error, info};
14use uuid::Uuid;
15
16use ironflow_store::models::{
17 NewRun, NewStep, RunStatus, RunUpdate, StepKind, StepStatus, StepUpdate, TriggerKind,
18 step_trace_id,
19};
20
21use crate::config::WorkflowStepConfig;
22use crate::context::WorkflowContext;
23use crate::error::EngineError;
24use crate::executor::SubWorkflowOutput;
25use crate::guard::WorkflowRejection;
26use crate::handler::{TypedWorkflow, WorkflowHandler};
27use crate::plan::{SharedPlanRecorder, lock_plan};
28
29impl WorkflowContext {
30 /// Execute a sub-workflow step.
31 ///
32 /// Creates a child run of `handler` whose payload is `input`, executes it
33 /// with its own steps and lifecycle, and returns its run ID and aggregated
34 /// metrics. The child declares its input type through [`TypedWorkflow`],
35 /// so only a `W::Input` is accepted.
36 ///
37 /// Requires the context to be created with
38 /// `with_handler_resolver`.
39 ///
40 /// # Errors
41 ///
42 /// Returns [`EngineError::InvalidWorkflow`] if no handler is registered
43 /// with the given name, or if no handler resolver is available, and
44 /// [`EngineError::Serialization`] if `input` cannot be serialized.
45 ///
46 /// # Examples
47 ///
48 /// ```no_run
49 /// use ironflow_engine::context::WorkflowContext;
50 /// use ironflow_engine::error::EngineError;
51 /// use ironflow_engine::handler::{HandlerFuture, TypedWorkflow, WorkflowHandler};
52 /// use serde::{Deserialize, Serialize};
53 ///
54 /// #[derive(Serialize, Deserialize)]
55 /// struct CollectInput {
56 /// scope: String,
57 /// }
58 ///
59 /// struct Collect;
60 ///
61 /// impl WorkflowHandler for Collect {
62 /// fn name(&self) -> &str { "collect" }
63 /// fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
64 /// Box::pin(async move { Ok(()) })
65 /// }
66 /// }
67 ///
68 /// impl TypedWorkflow for Collect {
69 /// type Input = CollectInput;
70 /// }
71 ///
72 /// # async fn example(ctx: &mut WorkflowContext) -> Result<(), EngineError> {
73 /// let child = ctx.workflow(&Collect, CollectInput { scope: "system".to_string() }).await?;
74 /// let steps = ctx.store().list_steps(child.run_id()).await?;
75 /// # Ok(())
76 /// # }
77 /// ```
78 ///
79 /// Any other input type is a compile error:
80 ///
81 /// ```compile_fail,E0308
82 /// # use ironflow_engine::context::WorkflowContext;
83 /// # use ironflow_engine::error::EngineError;
84 /// # use ironflow_engine::handler::{HandlerFuture, TypedWorkflow, WorkflowHandler};
85 /// # #[derive(serde::Serialize, serde::Deserialize)]
86 /// # struct CollectInput { scope: String }
87 /// # struct Collect;
88 /// # impl WorkflowHandler for Collect {
89 /// # fn name(&self) -> &str { "collect" }
90 /// # fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
91 /// # Box::pin(async move { Ok(()) })
92 /// # }
93 /// # }
94 /// # impl TypedWorkflow for Collect { type Input = CollectInput; }
95 /// # async fn example(ctx: &mut WorkflowContext) -> Result<(), EngineError> {
96 /// ctx.workflow(&Collect, serde_json::json!({"scope": "system"})).await?;
97 /// # Ok(())
98 /// # }
99 /// ```
100 pub async fn workflow<W: TypedWorkflow>(
101 &mut self,
102 handler: &W,
103 input: W::Input,
104 ) -> Result<SubWorkflowOutput, EngineError> {
105 let payload = to_value(&input)?;
106 self.run_sub_workflow(handler, payload).await
107 }
108
109 /// Execute a sub-workflow step whose child is only known at run time.
110 ///
111 /// Same as [`workflow`](Self::workflow), without the compile-time check of
112 /// the payload: the child must deserialize `payload` itself.
113 ///
114 /// # Errors
115 ///
116 /// Same as [`workflow`](Self::workflow).
117 ///
118 /// # Examples
119 ///
120 /// ```no_run
121 /// use ironflow_engine::context::WorkflowContext;
122 /// use ironflow_engine::error::EngineError;
123 /// use ironflow_engine::handler::WorkflowHandler;
124 /// use serde_json::json;
125 ///
126 /// # #[allow(deprecated)]
127 /// # async fn example(ctx: &mut WorkflowContext, child: &dyn WorkflowHandler) -> Result<(), EngineError> {
128 /// let result = ctx.workflow_dyn(child, json!({"scope": "system"})).await?;
129 /// println!("child run {}", result.run_id());
130 /// # Ok(())
131 /// # }
132 /// ```
133 #[deprecated(
134 note = "implement `TypedWorkflow` on the child and call `workflow`: its payload is then checked at compile time"
135 )]
136 pub async fn workflow_dyn(
137 &mut self,
138 handler: &dyn WorkflowHandler,
139 payload: Value,
140 ) -> Result<SubWorkflowOutput, EngineError> {
141 self.run_sub_workflow(handler, payload).await
142 }
143
144 /// Record, then run or plan, a sub-workflow step.
145 async fn run_sub_workflow(
146 &mut self,
147 handler: &dyn WorkflowHandler,
148 payload: Value,
149 ) -> Result<SubWorkflowOutput, EngineError> {
150 // Plan mode: record the invocation, expand the child handler in the
151 // same recorder, and return a synthetic output. No child run is
152 // created and no step of the child is executed.
153 if let Some(plan) = self.plan().cloned() {
154 return self.plan_sub_workflow(&plan, handler, payload).await;
155 }
156
157 // Guard check: verify limits before creating the step.
158 if let (Some(guard_config), Some(guard_state)) = (&self.guard_config, &self.guard_state) {
159 let state = guard_state
160 .lock()
161 .map_err(|_| WorkflowRejection::GuardUnavailable)?;
162 state.check(guard_config, handler.name())?;
163 }
164
165 let config = WorkflowStepConfig::new(handler.name(), payload);
166 let position = self.position;
167 self.position += 1;
168
169 let trace_id = step_trace_id(self.run_id, &config.workflow_name, position);
170 let step = self
171 .store
172 .create_step(NewStep {
173 run_id: self.run_id,
174 trace_id,
175 name: config.workflow_name.clone(),
176 kind: StepKind::Workflow,
177 position,
178 input: Some(to_value(&config)?),
179 is_error_handler: false,
180 })
181 .await?;
182
183 self.start_step(step.id, Utc::now()).await?;
184
185 // Record invocation in guard state (fail-closed).
186 if let Some(guard_state) = &self.guard_state {
187 let mut state = guard_state
188 .lock()
189 .map_err(|_| WorkflowRejection::GuardUnavailable)?;
190 state.record_invocation(handler.name());
191 }
192
193 match self.execute_child_workflow(&config).await {
194 Ok((output, child_had_allowed_failure)) => {
195 self.total_cost_usd += output.cost_usd();
196 self.total_duration_ms += output.duration_ms();
197 if child_had_allowed_failure {
198 self.has_allowed_failure = true;
199 }
200
201 let completed_at = Utc::now();
202 self.store
203 .update_step(
204 step.id,
205 StepUpdate {
206 status: Some(StepStatus::Completed),
207 output: Some(to_value(&output)?),
208 duration_ms: Some(output.duration_ms()),
209 cost_usd: Some(output.cost_usd()),
210 completed_at: Some(completed_at),
211 ..StepUpdate::default()
212 },
213 )
214 .await?;
215
216 info!(
217 run_id = %self.run_id,
218 child_workflow = %config.workflow_name,
219 duration_ms = output.duration_ms(),
220 "workflow step completed"
221 );
222
223 self.last_step_ids = vec![step.id];
224
225 self.guard_record_return();
226 Ok(output)
227 }
228 Err(err) => {
229 let completed_at = Utc::now();
230 if let Err(store_err) = self
231 .store
232 .update_step(
233 step.id,
234 StepUpdate {
235 status: Some(StepStatus::Failed),
236 error: Some(err.to_string()),
237 completed_at: Some(completed_at),
238 ..StepUpdate::default()
239 },
240 )
241 .await
242 {
243 error!(step_id = %step.id, error = %store_err, "failed to persist step failure");
244 }
245
246 self.guard_record_return();
247 Err(err)
248 }
249 }
250 }
251
252 /// Record a sub-workflow invocation while planning, expanding the child
253 /// handler into the same plan when the depth limit allows it.
254 ///
255 /// The child plans against its own payload and under its own workflow
256 /// name; the parent's payload is restored on the way out.
257 async fn plan_sub_workflow(
258 &mut self,
259 plan: &SharedPlanRecorder,
260 handler: &dyn WorkflowHandler,
261 payload: Value,
262 ) -> Result<SubWorkflowOutput, EngineError> {
263 self.position += 1;
264 let sub_name = handler.name().to_string();
265 // No child run exists while planning: a nil id and zero metrics.
266 let planned = SubWorkflowOutput::new(
267 Uuid::nil(),
268 &sub_name,
269 RunStatus::Completed,
270 Decimal::ZERO,
271 0,
272 );
273
274 {
275 let mut recorder = lock_plan(plan);
276 if !recorder.record(&sub_name, StepKind::Workflow, &self.workflow_name, None) {
277 return Ok(planned);
278 }
279 recorder.set_last(vec![sub_name.clone()]);
280 }
281
282 let expand = lock_plan(plan).enter_workflow();
283 if expand {
284 let previous_payload = lock_plan(plan).swap_payload(payload.clone());
285
286 let mut child = WorkflowContext::new(
287 Uuid::now_v7(),
288 sub_name.clone(),
289 self.store.clone(),
290 self.provider.clone(),
291 );
292 child.handler_resolver = self.handler_resolver.clone();
293 child.set_plan(plan.clone());
294
295 if let Err(err) = handler.execute(&mut child).await {
296 lock_plan(plan).fail(format!(
297 "sub-workflow {sub_name} could not be planned: {err}"
298 ));
299 }
300
301 let mut recorder = lock_plan(plan);
302 recorder.swap_payload(previous_payload);
303 recorder.leave_workflow();
304 }
305
306 Ok(planned)
307 }
308
309 /// Execute a child workflow and return aggregated output plus whether
310 /// at least one `allow_failure` step failed.
311 async fn execute_child_workflow(
312 &self,
313 config: &WorkflowStepConfig,
314 ) -> Result<(SubWorkflowOutput, bool), EngineError> {
315 let resolver = self.handler_resolver.as_ref().ok_or_else(|| {
316 EngineError::InvalidWorkflow(
317 "sub-workflow requires a handler resolver (use Engine to execute)".to_string(),
318 )
319 })?;
320
321 let handler = resolver(&config.workflow_name).ok_or_else(|| {
322 EngineError::InvalidWorkflow(format!("no handler registered: {}", config.workflow_name))
323 })?;
324
325 // A child run inherits both the parent labels and the parent author:
326 // whoever triggered the parent workflow is accountable for its children.
327 let parent = self.store.get_run(self.run_id).await?;
328 let (parent_labels, parent_author) =
329 parent.map(|r| (r.labels, r.created_by)).unwrap_or_default();
330
331 let child_run = self
332 .store
333 .create_run(NewRun {
334 workflow_name: config.workflow_name.clone(),
335 trigger: TriggerKind::Workflow,
336 payload: config.payload.clone(),
337 max_retries: 0,
338 handler_version: None,
339 labels: parent_labels,
340 scheduled_at: None,
341 created_by: parent_author,
342 idempotency_key: None,
343 // The child shares the parent's cap; it does not get its own budget.
344 max_cost_usd: self.max_cost_usd,
345 })
346 .await?
347 .into_run();
348
349 let child_run_id = child_run.id;
350 info!(
351 parent_run_id = %self.run_id,
352 child_run_id = %child_run_id,
353 workflow = %config.workflow_name,
354 "child run created"
355 );
356
357 self.store
358 .update_run_status(child_run_id, RunStatus::Running)
359 .await?;
360
361 let run_start = Instant::now();
362 let mut child_ctx = WorkflowContext {
363 run_id: child_run_id,
364 workflow_name: config.workflow_name.clone(),
365 store: self.store.clone(),
366 provider: self.provider.clone(),
367 decision_provider: self.decision_provider.clone(),
368 handler_resolver: self.handler_resolver.clone(),
369 position: 0,
370 last_step_ids: Vec::new(),
371 total_cost_usd: Decimal::ZERO,
372 total_duration_ms: 0,
373 max_cost_usd: self.max_cost_usd,
374 // Everything the parent chain already spent counts against the
375 // shared cap, so the child cannot restart the budget from zero.
376 inherited_cost_usd: self.charged_cost_usd(),
377 replay_steps: HashMap::new(),
378 replay_wave_steps: HashMap::new(),
379 granted_approvals: HashMap::new(),
380 answered_inputs: HashMap::new(),
381 // A child run is created fresh here; it is never itself retried.
382 attempt: 1,
383 carried_duration_ms: 0,
384 log_sender: self.log_sender.clone(),
385 // A child shares the storage backend but not the parent's artifacts:
386 // input lookups are scoped to the child's own run.
387 artifact_sink: self.artifact_sink.clone(),
388 has_allowed_failure: false,
389 error_handlers: Vec::new(),
390 guard_state: self.guard_state.clone(),
391 guard_config: self.guard_config.clone(),
392 step_results: Vec::new(),
393 event_bus: self.event_bus.clone(),
394 // A child run is mocked exactly like its parent.
395 interceptor: self.interceptor.clone(),
396 trace_context: self.trace_context.child(),
397 operation_ctx: None,
398 plan: None,
399 };
400
401 let result = handler.execute(&mut child_ctx).await;
402 let total_duration = run_start.elapsed().as_millis() as u64;
403 let completed_at = Utc::now();
404
405 match result {
406 Ok(()) => {
407 let child_status = if child_ctx.has_allowed_failure {
408 RunStatus::Warning
409 } else {
410 RunStatus::Completed
411 };
412 self.store
413 .update_run(
414 child_run_id,
415 RunUpdate {
416 status: Some(child_status),
417 cost_usd: Some(child_ctx.total_cost_usd),
418 duration_ms: Some(total_duration),
419 completed_at: Some(completed_at),
420 ..RunUpdate::default()
421 },
422 )
423 .await?;
424
425 let child_had_allowed_failure = child_ctx.has_allowed_failure;
426 Ok((
427 SubWorkflowOutput::new(
428 child_run_id,
429 &config.workflow_name,
430 child_status,
431 child_ctx.total_cost_usd,
432 total_duration,
433 ),
434 child_had_allowed_failure,
435 ))
436 }
437 Err(err) => {
438 if let Err(store_err) = self
439 .store
440 .update_run(
441 child_run_id,
442 RunUpdate {
443 status: Some(RunStatus::Failed),
444 error: Some(err.to_string()),
445 cost_usd: Some(child_ctx.total_cost_usd),
446 duration_ms: Some(total_duration),
447 completed_at: Some(completed_at),
448 ..RunUpdate::default()
449 },
450 )
451 .await
452 {
453 error!(
454 child_run_id = %child_run_id,
455 store_error = %store_err,
456 "failed to persist child run failure"
457 );
458 }
459
460 Err(err)
461 }
462 }
463 }
464}