Skip to main content

beam_core/workflow_runtime/
dispatch.rs

1use 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                    // --- prepare (parse + canonicalise) BEFORE any side-effect ---
178                    let prepared = hooks
179                        .prepare_host_executor(&executor.executor, &resolved_input)
180                        .context("prepare_host_executor failed")?;
181
182                    // --- write effect-input.json using the canonical input ---
183                    let _ = write_effect_input_sidecar(
184                        &rt.log,
185                        activity_id,
186                        &attempt_id,
187                        &prepared.canonical_input,
188                    )
189                    .await?;
190
191                    // --- emit effectAttempted BEFORE calling the external provider ---
192                    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}