beam_core/workflow_definition/
schema.rs1use std::collections::BTreeMap;
2
3use serde::{Deserialize, Serialize};
4use serde_json::Value;
5
6#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
7#[serde(rename_all = "camelCase")]
8pub struct ParamDef {
9 #[serde(rename = "type")]
10 pub param_type: String,
11 #[serde(default, skip_serializing_if = "Option::is_none")]
12 pub format: Option<String>,
13 #[serde(default, skip_serializing_if = "Option::is_none")]
14 pub required: Option<bool>,
15 #[serde(default, skip_serializing_if = "Option::is_none")]
16 pub default: Option<Value>,
17 #[serde(default, skip_serializing_if = "Option::is_none")]
18 pub description: Option<String>,
19}
20
21#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
22#[serde(rename_all = "camelCase")]
23pub struct RetryPolicy {
24 pub max_attempts: u64,
25 pub backoff: String,
26 pub base_ms: u64,
27 #[serde(default, skip_serializing_if = "Option::is_none")]
28 pub factor: Option<f64>,
29 #[serde(default, skip_serializing_if = "Option::is_none")]
30 pub jitter: Option<bool>,
31}
32
33#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
34#[serde(rename_all = "camelCase")]
35pub struct HumanGate {
36 pub stage: String,
37 pub prompt: Value,
38 #[serde(default, skip_serializing_if = "Option::is_none")]
39 pub approvers: Option<Vec<String>>,
40 #[serde(default, skip_serializing_if = "Option::is_none")]
41 pub deadline_ms: Option<u64>,
42 #[serde(default, skip_serializing_if = "Option::is_none")]
43 pub on_timeout: Option<String>,
44}
45
46#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
47#[serde(rename_all = "camelCase")]
48pub struct LoopTerminate {
49 pub node: String,
50 pub via: String,
51}
52
53#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
54#[serde(rename_all = "camelCase")]
55pub struct LoopOutputProjection {
56 pub from: String,
57}
58
59#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
60#[serde(rename_all = "camelCase")]
61pub struct NodeBase {
62 #[serde(default, skip_serializing_if = "Option::is_none")]
63 pub description: Option<String>,
64 #[serde(default, skip_serializing_if = "Option::is_none")]
65 pub depends: Option<Vec<String>>,
66 #[serde(default, skip_serializing_if = "Option::is_none")]
67 pub human_gate: Option<HumanGate>,
68 #[serde(default, skip_serializing_if = "Option::is_none")]
69 pub retry_policy: Option<RetryPolicy>,
70 #[serde(default, skip_serializing_if = "Option::is_none")]
71 pub timeout_ms: Option<u64>,
72 #[serde(default, skip_serializing_if = "Option::is_none")]
73 pub max_output_bytes: Option<u64>,
74 #[serde(default, skip_serializing_if = "Option::is_none")]
75 pub output_schema: Option<Value>,
76 #[serde(default, skip_serializing_if = "Option::is_none")]
77 pub unsafe_allow_ungated: Option<bool>,
78}
79
80#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
81#[serde(rename_all = "camelCase")]
82pub struct SubagentNode {
83 #[serde(flatten)]
84 pub base: NodeBase,
85 pub bot: String,
86 pub prompt: Value,
87 #[serde(default, skip_serializing_if = "Option::is_none")]
88 pub working_dir: Option<String>,
89 #[serde(default, skip_serializing_if = "Option::is_none")]
90 pub model_overrides: Option<Value>,
91 #[serde(default, skip_serializing_if = "Option::is_none")]
92 pub tool_policy: Option<Value>,
93}
94
95#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
96#[serde(rename_all = "camelCase")]
97pub struct HostExecutorNode {
98 #[serde(flatten)]
99 pub base: NodeBase,
100 pub executor: String,
101 pub input: Value,
102}
103
104#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
105#[serde(rename_all = "camelCase")]
106pub struct LoopNode {
107 #[serde(flatten)]
108 pub base: NodeBase,
109 pub max_iterations: u64,
110 pub body: Vec<String>,
111 pub terminate: LoopTerminate,
112 #[serde(default, skip_serializing_if = "Option::is_none")]
113 pub output: Option<LoopOutputProjection>,
114}
115
116#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
117#[serde(rename_all = "camelCase")]
118pub struct DecisionNode {
119 #[serde(flatten)]
120 pub base: NodeBase,
121}
122
123#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
124#[serde(tag = "type", rename_all = "camelCase")]
125pub enum WorkflowNode {
126 Subagent(SubagentNode),
127 HostExecutor(HostExecutorNode),
128 Loop(LoopNode),
129 Decision(DecisionNode),
130}
131
132#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
133#[serde(rename_all = "camelCase")]
134pub struct WorkflowDefaults {
135 #[serde(default, skip_serializing_if = "Option::is_none")]
136 pub retry_policy: Option<RetryPolicy>,
137 #[serde(default, skip_serializing_if = "Option::is_none")]
138 pub timeout_ms: Option<u64>,
139 #[serde(default, skip_serializing_if = "Option::is_none")]
140 pub max_output_bytes: Option<u64>,
141 #[serde(default, skip_serializing_if = "Option::is_none")]
142 pub max_concurrency: Option<u64>,
143}
144
145#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
146#[serde(rename_all = "camelCase")]
147pub struct WorkflowDefinition {
148 pub workflow_id: String,
149 pub version: u64,
150 #[serde(default, skip_serializing_if = "Option::is_none")]
151 pub params: Option<BTreeMap<String, ParamDef>>,
152 #[serde(default, skip_serializing_if = "Option::is_none")]
153 pub defaults: Option<WorkflowDefaults>,
154 pub nodes: BTreeMap<String, WorkflowNode>,
155}