beam_core/workflow_runtime/
dispatch.rs1use anyhow::{Context, Result};
2use serde_json::Value;
3
4use crate::workflow_binding::{BindingContext, resolve_bindings, resolve_bound_string};
5use crate::workflow_orchestrator::OrchestratorAction;
6use crate::workflow_sidecar::write_effect_input_sidecar;
7use crate::{EventDraft, RunSnapshotDTO, WorkflowActor, WorkflowNode, read_run_snapshot};
8
9use super::completion::settle_work_result;
10use super::helpers::{
11 derive_workflow_idempotency_key, gate_attempt_id, loop_context_from_activity, now_ms,
12 sha256_hex, split_prompt, work_attempt_id, write_json_blob,
13};
14use super::{WorkflowDispatchOutcome, WorkflowDispatchRun, WorkflowRuntimeContext};
15
16async fn read_snapshot(rt: &WorkflowRuntimeContext) -> Result<RunSnapshotDTO> {
17 read_run_snapshot(&rt.log.run_dir)
18 .await?
19 .context("workflow runtime requires an existing run snapshot")
20}
21
22pub async fn dispatch_gate(
23 rt: &mut WorkflowRuntimeContext,
24 action: &crate::OrchestratorAction,
25) -> Result<()> {
26 match action {
27 OrchestratorAction::DispatchGate {
28 node_id,
29 activity_id,
30 human_gate,
31 } => {
32 let attempt_id = gate_attempt_id(activity_id);
33 let input_ref = write_json_blob(
34 &mut rt.log,
35 serde_json::json!({
36 "kind": "human-gate",
37 "prompt": human_gate.prompt,
38 "approvers": human_gate.approvers,
39 }),
40 )?;
41 rt.log.append(EventDraft {
42 event_type: "attemptCreated".to_string(),
43 actor: WorkflowActor::Scheduler,
44 payload: serde_json::json!({
45 "nodeId": node_id,
46 "activityId": activity_id,
47 "attemptId": attempt_id,
48 "attemptNumber": 1,
49 "inputRef": input_ref,
50 }),
51 timestamp: None,
52 payload_hash: None,
53 })?;
54
55 let snap = read_snapshot(rt).await?;
56 let ctx = BindingContext {
57 snapshot: &snap,
58 def: &rt.def,
59 run_dir: &rt.log.run_dir,
60 loop_context: loop_context_from_activity(activity_id),
61 };
62 let resolved_prompt = resolve_bound_string(&human_gate.prompt, &ctx).await?;
63 let prompt_field = split_prompt(&mut rt.log, &resolved_prompt)?;
64 let _ = crate::create_wait(
65 &mut rt.log,
66 crate::CreateWaitInput {
67 activity_id: activity_id.clone(),
68 attempt_id,
69 node_id: node_id.clone(),
70 wait_kind: crate::WaitKind::HumanGate,
71 deadline_at: human_gate.deadline_ms.map(|ms| now_ms() + ms),
72 prompt: prompt_field.prompt,
73 prompt_ref: prompt_field.prompt_ref,
74 prompt_preview: prompt_field.prompt_preview,
75 approvers: human_gate.approvers.clone(),
76 on_timeout: human_gate.on_timeout.as_deref().map(|v| match v {
77 "success" => crate::WaitOnTimeout::Success,
78 _ => crate::WaitOnTimeout::Fail,
79 }),
80 },
81 )
82 .await?;
83 Ok(())
84 }
85 _ => anyhow::bail!("dispatch_gate called with non-gate action"),
86 }
87}
88
89pub async fn dispatch_work<H: super::WorkflowExecutionHooks>(
90 rt: &mut WorkflowRuntimeContext,
91 hooks: &mut H,
92 action: &crate::OrchestratorAction,
93) -> Result<WorkflowDispatchOutcome> {
94 match action {
95 OrchestratorAction::DispatchWork {
96 node_id,
97 activity_id,
98 node,
99 } => {
100 let attempt_id = work_attempt_id(activity_id, 1);
101 let input_ref = write_json_blob(
102 &mut rt.log,
103 serde_json::json!({
104 "kind": match node {
105 WorkflowNode::Subagent(_) => "subagent",
106 WorkflowNode::HostExecutor(_) => "hostExecutor",
107 WorkflowNode::Loop(_) => "loop",
108 WorkflowNode::Decision(_) => "decision",
109 },
110 "bot_or_executor": match node {
111 WorkflowNode::Subagent(n) => Value::String(n.bot.clone()),
112 WorkflowNode::HostExecutor(n) => Value::String(n.executor.clone()),
113 _ => Value::Null,
114 },
115 "prompt_or_input": match node {
116 WorkflowNode::Subagent(n) => n.prompt.clone(),
117 WorkflowNode::HostExecutor(n) => n.input.clone(),
118 _ => Value::Null,
119 }
120 }),
121 )?;
122 rt.log.append(EventDraft {
123 event_type: "attemptCreated".to_string(),
124 actor: WorkflowActor::Scheduler,
125 payload: serde_json::json!({
126 "nodeId": node_id,
127 "activityId": activity_id,
128 "attemptId": attempt_id,
129 "attemptNumber": 1,
130 "inputRef": input_ref,
131 }),
132 timestamp: None,
133 payload_hash: None,
134 })?;
135
136 let snap = read_snapshot(rt).await?;
137 let bind_ctx = BindingContext {
138 snapshot: &snap,
139 def: &rt.def,
140 run_dir: &rt.log.run_dir,
141 loop_context: loop_context_from_activity(activity_id),
142 };
143
144 match node {
145 WorkflowNode::Subagent(subagent) => {
146 let resolved_prompt = resolve_bound_string(&subagent.prompt, &bind_ctx).await?;
147 rt.log.append(EventDraft {
148 event_type: "activityRunning".to_string(),
149 actor: WorkflowActor::Scheduler,
150 payload: serde_json::json!({
151 "activityId": activity_id,
152 "attemptId": attempt_id,
153 "leaseId": format!("lease-{}", attempt_id),
154 }),
155 timestamp: None,
156 payload_hash: None,
157 })?;
158 let result = hooks
159 .execute_subagent(
160 WorkflowDispatchRun {
161 run_id: &rt.log.run_id,
162 workflow_id: snap.run.workflow_id.as_deref().unwrap_or(""),
163 revision_id: snap.run.revision_id.as_deref().unwrap_or(""),
164 activity_id,
165 attempt_id: &attempt_id,
166 node_id,
167 },
168 subagent,
169 resolved_prompt,
170 )
171 .await?;
172 settle_work_result(&mut rt.log, activity_id, &attempt_id, result).await
173 }
174 WorkflowNode::HostExecutor(executor) => {
175 let resolved_input = resolve_bindings(&executor.input, &bind_ctx).await?;
176
177 let prepared = hooks
179 .prepare_host_executor(&executor.executor, &resolved_input)
180 .context("prepare_host_executor failed")?;
181
182 let _ = write_effect_input_sidecar(
184 &rt.log,
185 activity_id,
186 &attempt_id,
187 &prepared.canonical_input,
188 )
189 .await?;
190
191 let idempotency_key = derive_workflow_idempotency_key(
193 snap.run.workflow_id.as_deref().unwrap_or(""),
194 snap.run.revision_id.as_deref().unwrap_or(""),
195 &rt.log.run_id,
196 node_id,
197 &attempt_id,
198 );
199 let input_bytes = serde_json::to_vec(&prepared.canonical_input)?;
200 let input_hash = sha256_hex(&input_bytes);
201 rt.log.append(EventDraft {
202 event_type: "effectAttempted".to_string(),
203 actor: WorkflowActor::Scheduler,
204 payload: serde_json::json!({
205 "activityId": activity_id,
206 "attemptId": attempt_id,
207 "idempotencyKey": idempotency_key,
208 "inputHash": input_hash,
209 "idempotencyTtlMs": prepared.idempotency_ttl_ms,
210 "provider": prepared.provider,
211 }),
212 timestamp: None,
213 payload_hash: None,
214 })?;
215
216 let result = hooks
217 .execute_host_executor(
218 WorkflowDispatchRun {
219 run_id: &rt.log.run_id,
220 workflow_id: snap.run.workflow_id.as_deref().unwrap_or(""),
221 revision_id: snap.run.revision_id.as_deref().unwrap_or(""),
222 activity_id,
223 attempt_id: &attempt_id,
224 node_id,
225 },
226 executor,
227 prepared.parsed_input,
228 )
229 .await?;
230 settle_work_result(&mut rt.log, activity_id, &attempt_id, result).await
231 }
232 WorkflowNode::Loop(_) | WorkflowNode::Decision(_) => {
233 anyhow::bail!("dispatch_work received unsupported node type")
234 }
235 }
236 }
237 _ => anyhow::bail!("dispatch_work called with non-work action"),
238 }
239}