Skip to main content

beam_core/workflow_snapshot/
model.rs

1use std::collections::BTreeMap;
2
3use serde::{Deserialize, Serialize};
4use serde_json::Value;
5
6use crate::RunChatBinding;
7use crate::WorkflowOutputRef;
8
9#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
10#[serde(rename_all = "camelCase")]
11pub struct BlobPreviewDTO {
12    #[serde(default, skip_serializing_if = "Option::is_none")]
13    pub output_hash: Option<String>,
14    #[serde(default, skip_serializing_if = "Option::is_none")]
15    pub output_bytes: Option<usize>,
16    #[serde(default, skip_serializing_if = "Option::is_none")]
17    pub content_type: Option<String>,
18    #[serde(default, skip_serializing_if = "Option::is_none")]
19    pub truncated: Option<bool>,
20    #[serde(default, skip_serializing_if = "Option::is_none")]
21    pub value: Option<Value>,
22    #[serde(default, skip_serializing_if = "Option::is_none")]
23    pub text: Option<String>,
24    #[serde(default, skip_serializing_if = "Option::is_none")]
25    pub error: Option<String>,
26    #[serde(default, skip_serializing_if = "Option::is_none")]
27    pub redacted: Option<bool>,
28}
29
30#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
31#[serde(rename_all = "camelCase")]
32pub struct AttemptTerminalDTO {
33    pub session_id: String,
34    #[serde(default, skip_serializing_if = "Option::is_none")]
35    pub cli_session_id: Option<String>,
36    pub web_port: u16,
37    pub status: String,
38    #[serde(default, skip_serializing_if = "Option::is_none")]
39    pub lark_app_id: Option<String>,
40    #[serde(default, skip_serializing_if = "Option::is_none")]
41    pub bot_name: Option<String>,
42    #[serde(default, skip_serializing_if = "Option::is_none")]
43    pub cli_id: Option<String>,
44    #[serde(default, skip_serializing_if = "Option::is_none")]
45    pub working_dir: Option<String>,
46    #[serde(default, skip_serializing_if = "Option::is_none")]
47    pub log_path: Option<String>,
48    pub started_at: u64,
49    pub updated_at: u64,
50    #[serde(default, skip_serializing_if = "Option::is_none")]
51    pub closed_at: Option<u64>,
52    #[serde(default, skip_serializing_if = "Option::is_none")]
53    pub error: Option<String>,
54    #[serde(default, skip_serializing_if = "Option::is_none")]
55    pub has_terminal_log: Option<bool>,
56}
57
58#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
59#[serde(rename_all = "camelCase")]
60pub struct AttemptIODTO {
61    #[serde(default, skip_serializing_if = "Option::is_none")]
62    pub input: Option<BlobPreviewDTO>,
63    #[serde(default, skip_serializing_if = "Option::is_none")]
64    pub resolved_input: Option<BlobPreviewDTO>,
65    #[serde(default, skip_serializing_if = "Option::is_none")]
66    pub output: Option<BlobPreviewDTO>,
67    #[serde(default, skip_serializing_if = "Option::is_none")]
68    pub log: Option<BlobPreviewDTO>,
69    #[serde(default, skip_serializing_if = "Option::is_none")]
70    pub terminal: Option<AttemptTerminalDTO>,
71    #[serde(default, skip_serializing_if = "Option::is_none")]
72    pub wait_prompt: Option<BlobPreviewDTO>,
73}
74
75#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
76#[serde(rename_all = "camelCase")]
77pub enum RunStatus {
78    Pending,
79    Running,
80    Waiting,
81    Succeeded,
82    Failed,
83    Cancelled,
84}
85
86#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
87#[serde(rename_all = "camelCase")]
88pub enum NodeStatus {
89    Idle,
90    Triggered,
91    Running,
92    Waiting,
93    Retrying,
94    Succeeded,
95    Failed,
96    Skipped,
97    Cancelled,
98}
99
100#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
101#[serde(rename_all = "camelCase")]
102pub enum ActivityStatus {
103    Pending,
104    Acquired,
105    Running,
106    Waiting,
107    EffectAttempting,
108    Succeeded,
109    Failed,
110    TimedOut,
111    Cancelled,
112}
113
114#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
115#[serde(rename_all = "camelCase")]
116pub enum LoopIterationStatus {
117    Running,
118    Approved,
119    Rejected,
120    Failed,
121    Cancelled,
122}
123
124#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
125#[serde(rename_all = "camelCase")]
126pub enum LoopStatus {
127    Running,
128    Succeeded,
129    Failed,
130    Cancelled,
131}
132
133#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
134#[serde(rename_all = "camelCase")]
135pub struct EffectAttemptedState {
136    pub idempotency_key: String,
137    pub input_hash: String,
138    pub idempotency_ttl_ms: u64,
139    pub provider: String,
140    pub attempted_at_event_id: String,
141    pub attempted_at_ms: u64,
142}
143
144#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
145#[serde(rename_all = "camelCase")]
146pub struct ReconcileResultState {
147    pub decision: String,
148    pub capability: String,
149    pub evidence: Value,
150    pub event_id: String,
151}
152
153#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
154#[serde(rename_all = "camelCase")]
155pub struct WaitResolutionState {
156    pub kind: String,
157    #[serde(default, skip_serializing_if = "Option::is_none")]
158    pub resolution: Option<String>,
159    #[serde(default, skip_serializing_if = "Option::is_none")]
160    pub by: Option<String>,
161    #[serde(default, skip_serializing_if = "Option::is_none")]
162    pub comment: Option<String>,
163    #[serde(default, skip_serializing_if = "Option::is_none")]
164    pub event_id: Option<String>,
165    #[serde(default, skip_serializing_if = "Option::is_none")]
166    pub deadline_at: Option<u64>,
167    #[serde(default, skip_serializing_if = "Option::is_none")]
168    pub exceeded_at_ms: Option<u64>,
169}
170
171#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
172#[serde(rename_all = "camelCase")]
173pub struct WaitState {
174    pub wait_kind: String,
175    #[serde(default, skip_serializing_if = "Option::is_none")]
176    pub deadline_at: Option<u64>,
177    #[serde(default, skip_serializing_if = "Option::is_none")]
178    pub prompt: Option<String>,
179    #[serde(default, skip_serializing_if = "Option::is_none")]
180    pub prompt_ref: Option<WorkflowOutputRef>,
181    #[serde(default, skip_serializing_if = "Option::is_none")]
182    pub prompt_preview: Option<String>,
183    #[serde(default, skip_serializing_if = "Option::is_none")]
184    pub approvers: Option<Vec<String>>,
185    #[serde(default, skip_serializing_if = "Option::is_none")]
186    pub on_timeout: Option<String>,
187    #[serde(default, skip_serializing_if = "Option::is_none")]
188    pub resolution: Option<WaitResolutionState>,
189}
190
191#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
192#[serde(rename_all = "camelCase")]
193pub struct CancelRequestState {
194    pub cancel_origin_event_id: String,
195    pub requested_by: String,
196    pub reason: String,
197    pub delivered: bool,
198}
199
200#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
201#[serde(rename_all = "camelCase")]
202pub struct AttemptState {
203    pub attempt_id: String,
204    pub attempt_number: u64,
205    pub input_ref: WorkflowOutputRef,
206    pub status: ActivityStatus,
207    #[serde(default, skip_serializing_if = "Option::is_none")]
208    pub lease_id: Option<String>,
209    #[serde(default, skip_serializing_if = "Option::is_none")]
210    pub timeout_ms: Option<u64>,
211    #[serde(default, skip_serializing_if = "Option::is_none")]
212    pub max_output_bytes: Option<u64>,
213    #[serde(default, skip_serializing_if = "Option::is_none")]
214    pub effect_attempted: Option<EffectAttemptedState>,
215    #[serde(default, skip_serializing_if = "Option::is_none")]
216    pub latest_reconcile_result: Option<ReconcileResultState>,
217    #[serde(default, skip_serializing_if = "Option::is_none")]
218    pub cancel_request: Option<CancelRequestState>,
219    #[serde(default, skip_serializing_if = "Option::is_none")]
220    pub wait: Option<WaitState>,
221    #[serde(default, skip_serializing_if = "Option::is_none")]
222    pub output: Option<WorkflowOutputRef>,
223    #[serde(default, skip_serializing_if = "Option::is_none")]
224    pub external_refs: Option<Value>,
225    #[serde(default, skip_serializing_if = "Option::is_none")]
226    pub error: Option<Value>,
227    #[serde(default, skip_serializing_if = "Option::is_none")]
228    pub running_ms: Option<u64>,
229    #[serde(default, skip_serializing_if = "Option::is_none")]
230    pub cancel_origin_event_id: Option<String>,
231}
232
233#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
234#[serde(rename_all = "camelCase")]
235pub struct ActivityState {
236    pub activity_id: String,
237    pub attempts: Vec<AttemptState>,
238    pub status: ActivityStatus,
239    #[serde(default, skip_serializing_if = "Option::is_none")]
240    pub current_attempt_id: Option<String>,
241    #[serde(default, skip_serializing_if = "Option::is_none")]
242    pub owner_node_id: Option<String>,
243}
244
245#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
246#[serde(rename_all = "camelCase")]
247pub struct NodeState {
248    pub node_id: String,
249    pub status: NodeStatus,
250    #[serde(default, skip_serializing_if = "Option::is_none")]
251    pub activity_id: Option<String>,
252    pub retry_count: u64,
253    #[serde(default, skip_serializing_if = "Option::is_none")]
254    pub next_attempt_at: Option<u64>,
255    #[serde(default, skip_serializing_if = "Option::is_none")]
256    pub error_class: Option<String>,
257    #[serde(default, skip_serializing_if = "Option::is_none")]
258    pub condition_event_id: Option<String>,
259    #[serde(default, skip_serializing_if = "Option::is_none")]
260    pub cancel_origin_event_id: Option<String>,
261}
262
263#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
264#[serde(rename_all = "camelCase")]
265pub struct LoopIterationState {
266    pub iteration: u64,
267    pub status: LoopIterationStatus,
268    pub body_activity_ids: Vec<String>,
269    #[serde(default, skip_serializing_if = "Option::is_none")]
270    pub decision_activity_id: Option<String>,
271    #[serde(default, skip_serializing_if = "Option::is_none")]
272    pub wait_resolved_event_id: Option<String>,
273    #[serde(default, skip_serializing_if = "Option::is_none")]
274    pub decision_by: Option<String>,
275    #[serde(default, skip_serializing_if = "Option::is_none")]
276    pub decision_comment: Option<String>,
277    #[serde(default, skip_serializing_if = "Option::is_none")]
278    pub timed_out: Option<bool>,
279}
280
281#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
282#[serde(rename_all = "camelCase")]
283pub struct LoopState {
284    pub loop_id: String,
285    pub status: LoopStatus,
286    pub iteration: u64,
287    pub max_iterations: u64,
288    pub iterations: Vec<LoopIterationState>,
289    #[serde(default, skip_serializing_if = "Option::is_none")]
290    pub output: Option<WorkflowOutputRef>,
291    #[serde(default, skip_serializing_if = "Option::is_none")]
292    pub error_code: Option<String>,
293    #[serde(default, skip_serializing_if = "Option::is_none")]
294    pub error_class: Option<String>,
295}
296
297#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
298#[serde(rename_all = "camelCase")]
299pub struct RunState {
300    pub run_id: String,
301    pub status: RunStatus,
302    #[serde(default, skip_serializing_if = "Option::is_none")]
303    pub workflow_id: Option<String>,
304    #[serde(default, skip_serializing_if = "Option::is_none")]
305    pub revision_id: Option<String>,
306    #[serde(default, skip_serializing_if = "Option::is_none")]
307    pub initiator: Option<String>,
308    #[serde(default, skip_serializing_if = "Option::is_none")]
309    pub input: Option<WorkflowOutputRef>,
310    #[serde(default, skip_serializing_if = "Option::is_none")]
311    pub output: Option<WorkflowOutputRef>,
312    #[serde(default, skip_serializing_if = "Option::is_none")]
313    pub failed_node_id: Option<String>,
314    #[serde(default, skip_serializing_if = "Option::is_none")]
315    pub root_cause_event_id: Option<String>,
316    #[serde(default, skip_serializing_if = "Option::is_none")]
317    pub cancel_origin_event_id: Option<String>,
318    #[serde(default, skip_serializing_if = "Option::is_none")]
319    pub bot_snapshots: Option<BTreeMap<String, BotSnapshot>>,
320    #[serde(default, skip_serializing_if = "Option::is_none")]
321    pub cancelled_run_intent: Option<CancelIntent>,
322    #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
323    pub cancelled_node_intents: BTreeMap<String, CancelIntent>,
324}
325
326#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
327#[serde(rename_all = "camelCase")]
328pub struct CancelIntent {
329    pub cancel_origin_event_id: String,
330    pub requested_by: String,
331    pub reason: String,
332}
333
334#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
335#[serde(rename_all = "camelCase")]
336pub struct BotSnapshot {
337    #[serde(default, skip_serializing_if = "Option::is_none")]
338    pub lark_app_id: Option<String>,
339    #[serde(default, skip_serializing_if = "Option::is_none")]
340    pub cli_id: Option<String>,
341    #[serde(default, skip_serializing_if = "Option::is_none")]
342    pub display_name: Option<String>,
343    #[serde(default, skip_serializing_if = "Option::is_none")]
344    pub working_dir: Option<String>,
345}
346
347#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
348#[serde(rename_all = "camelCase")]
349pub struct LoopSnapshotDTO {
350    pub loop_id: String,
351    pub status: LoopStatus,
352    pub iteration: u64,
353    pub max_iterations: u64,
354    pub iterations: Vec<LoopIterationState>,
355    #[serde(default, skip_serializing_if = "Option::is_none")]
356    pub output: Option<WorkflowOutputRef>,
357    #[serde(default, skip_serializing_if = "Option::is_none")]
358    pub error_code: Option<String>,
359    #[serde(default, skip_serializing_if = "Option::is_none")]
360    pub error_class: Option<String>,
361}
362
363#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
364#[serde(rename_all = "camelCase")]
365pub struct RunSnapshotDTO {
366    pub run_id: String,
367    pub run: RunState,
368    pub last_seq: u64,
369    pub nodes: Vec<NodeState>,
370    pub activities: Vec<ActivityState>,
371    #[serde(default, skip_serializing_if = "Option::is_none")]
372    pub loops: Option<BTreeMap<String, LoopSnapshotDTO>>,
373    pub dangling: DanglingSnapshot,
374    pub outputs: BTreeMap<String, WorkflowOutputRef>,
375    pub attempt_io: BTreeMap<String, AttemptIODTO>,
376    #[serde(default, skip_serializing_if = "Option::is_none")]
377    pub chat_binding: Option<RunChatBinding>,
378    pub updated_at: u64,
379}
380
381#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
382#[serde(rename_all = "camelCase")]
383pub struct DanglingSnapshot {
384    pub activities: Vec<String>,
385    pub effect_attempted: Vec<String>,
386    pub waits: Vec<String>,
387    /// Activities where a wait resolution (waitResolved or waitDeadlineExceeded)
388    /// was written but the terminal activity event (activitySucceeded / activityFailed)
389    /// has not been written yet.  These should be materialised by recovery/resume.
390    pub wait_resolutions: Vec<String>,
391    pub cancels: Vec<String>,
392}