Skip to main content

beam_core/workflow_definition/
schema.rs

1use 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}