meerkat_workgraph/generated/
protocol_work_execution_flow_observation.rs1use crate::machines::work_execution_lifecycle::{
6 WorkExecutionLifecycleEffect, WorkExecutionLifecycleInput,
7 WorkExecutionLifecycleMachineAuthority, WorkExecutionLifecycleMachineMutator,
8 WorkExecutionLifecycleMachineTransition, WorkExecutionLifecycleMachineTransitionError,
9};
10
11#[derive(Debug, Clone)]
12pub struct WorkExecutionFlowObservationObligation {
13 pub binding_id: String,
14 pub run_id: String,
15}
16
17#[macro_export]
18macro_rules! work_execution_flow_observation_feedback_input_patterns {
19 () => {
20 $crate::machines::work_execution_lifecycle::WorkExecutionLifecycleInput::ObserveFlowRunning
21 | $crate::machines::work_execution_lifecycle::WorkExecutionLifecycleInput::ObserveFlowCompleted
22 | $crate::machines::work_execution_lifecycle::WorkExecutionLifecycleInput::ObserveFlowFailed { .. }
23 | $crate::machines::work_execution_lifecycle::WorkExecutionLifecycleInput::ObserveFlowCanceled { .. }
24 | $crate::machines::work_execution_lifecycle::WorkExecutionLifecycleInput::ObserveRunLost { .. }
25 };
26}
27
28pub fn extract_obligations(
29 transition: &WorkExecutionLifecycleMachineTransition,
30) -> Vec<WorkExecutionFlowObservationObligation> {
31 transition
32 .effects()
33 .iter()
34 .filter_map(|effect| match effect {
35 WorkExecutionLifecycleEffect::FlowLaunchAccepted { binding_id, run_id } => {
36 Some(WorkExecutionFlowObservationObligation {
37 binding_id: binding_id.clone(),
38 run_id: run_id.clone(),
39 })
40 }
41 _ => None,
42 })
43 .collect()
44}
45
46pub fn submit_observe_flow_running(
47 authority: &mut WorkExecutionLifecycleMachineAuthority,
48 _obligation: WorkExecutionFlowObservationObligation,
49) -> Result<WorkExecutionLifecycleMachineTransition, WorkExecutionLifecycleMachineTransitionError> {
50 let transition = authority.apply(WorkExecutionLifecycleInput::ObserveFlowRunning)?;
51 Ok(transition)
52}
53
54pub fn submit_observe_flow_completed(
55 authority: &mut WorkExecutionLifecycleMachineAuthority,
56 _obligation: WorkExecutionFlowObservationObligation,
57) -> Result<WorkExecutionLifecycleMachineTransition, WorkExecutionLifecycleMachineTransitionError> {
58 let transition = authority.apply(WorkExecutionLifecycleInput::ObserveFlowCompleted)?;
59 Ok(transition)
60}
61
62pub fn submit_observe_flow_failed(
63 authority: &mut WorkExecutionLifecycleMachineAuthority,
64 _obligation: WorkExecutionFlowObservationObligation,
65 observed_failure_detail: Option<String>,
66) -> Result<WorkExecutionLifecycleMachineTransition, WorkExecutionLifecycleMachineTransitionError> {
67 let transition = authority.apply(WorkExecutionLifecycleInput::ObserveFlowFailed {
68 detail: observed_failure_detail,
69 })?;
70 Ok(transition)
71}
72
73pub fn submit_observe_flow_canceled(
74 authority: &mut WorkExecutionLifecycleMachineAuthority,
75 _obligation: WorkExecutionFlowObservationObligation,
76 observed_cancellation_detail: Option<String>,
77) -> Result<WorkExecutionLifecycleMachineTransition, WorkExecutionLifecycleMachineTransitionError> {
78 let transition = authority.apply(WorkExecutionLifecycleInput::ObserveFlowCanceled {
79 detail: observed_cancellation_detail,
80 })?;
81 Ok(transition)
82}
83
84pub fn submit_observe_run_lost(
85 authority: &mut WorkExecutionLifecycleMachineAuthority,
86 _obligation: WorkExecutionFlowObservationObligation,
87 lost_run_detail: String,
88) -> Result<WorkExecutionLifecycleMachineTransition, WorkExecutionLifecycleMachineTransitionError> {
89 let transition = authority.apply(WorkExecutionLifecycleInput::ObserveRunLost {
90 detail: lost_run_detail,
91 })?;
92 Ok(transition)
93}