Skip to main content

backbone_integrations/application/workflows/
example_saga_workflow.rs

1//! ExampleSaga flow implementation
2//!
3//! Generated by metaphor-schema
4//!
5//! Replace with the business intent of your saga — what it creates/updates atomically,
6//! what it validates first, and what event it emits on success.
7
8use serde::{Deserialize, Serialize};
9use std::sync::Arc;
10use chrono;
11use backbone_core::flow::{WorkflowStep, WorkflowContext};
12
13/// Error type for flow execution
14#[derive(Debug, Clone, thiserror::Error)]
15pub enum FlowError {
16    #[error("No current step to execute")]
17    NoCurrentStep,
18
19    #[error("Step execution failed: {0}")]
20    StepFailed(String),
21
22    #[error("Condition evaluation failed: {0}")]
23    ConditionFailed(String),
24
25    #[error("Compensation failed: {0}")]
26    CompensationFailed(String),
27
28    #[error("Flow timed out")]
29    Timeout,
30
31    #[error("Flow cancelled")]
32    Cancelled,
33
34    #[error("Invalid state transition: {from} -> {to}")]
35    InvalidTransition { from: String, to: String },
36}
37
38/// Execution status for ExampleSaga flow
39#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
40#[serde(rename_all = "snake_case")]
41pub enum ExampleSagaFlowStatus {
42    /// Flow is pending execution
43    Pending,
44    /// Flow is currently running
45    Running,
46    /// Flow is waiting for an event or condition
47    Waiting,
48    /// Flow completed successfully
49    Completed,
50    /// Flow failed
51    Failed,
52    /// Flow was cancelled
53    Cancelled,
54    /// Flow is compensating (rolling back)
55    Compensating,
56}
57
58/// Steps in ExampleSaga flow
59#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
60#[serde(rename_all = "snake_case")]
61pub enum ExampleSagaFlowStep {
62    /// validate
63    Validate,
64    /// write_example
65    WriteExample,
66    /// done
67    Done,
68    /// rejected
69    Rejected,
70}
71
72/// Instance of ExampleSaga flow execution
73#[derive(Debug, Clone, Serialize, Deserialize)]
74pub struct ExampleSagaFlowInstance {
75    /// Unique instance ID
76    pub id: String,
77    /// Current status
78    pub status: ExampleSagaFlowStatus,
79    /// Current step
80    pub current_step: Option<ExampleSagaFlowStep>,
81    /// Flow context (variables)
82    pub context: serde_json::Value,
83    /// Completed steps
84    pub completed_steps: Vec<ExampleSagaFlowStep>,
85    /// Error if failed
86    pub error: Option<String>,
87    /// Created timestamp
88    pub created_at: chrono::DateTime<chrono::Utc>,
89    /// Updated timestamp
90    pub updated_at: chrono::DateTime<chrono::Utc>,
91}
92
93impl ExampleSagaFlowInstance {
94    /// Create a new flow instance
95    pub fn new(id: impl Into<String>) -> Self {
96        let now = chrono::Utc::now();
97        Self {
98            id: id.into(),
99            status: ExampleSagaFlowStatus::Pending,
100            current_step: None,
101            context: serde_json::json!({}),
102            completed_steps: Vec::new(),
103            error: None,
104            created_at: now,
105            updated_at: now,
106        }
107    }
108
109    /// Check if flow is complete
110    pub fn is_complete(&self) -> bool {
111        matches!(
112            self.status,
113            ExampleSagaFlowStatus::Completed | ExampleSagaFlowStatus::Failed | ExampleSagaFlowStatus::Cancelled
114        )
115    }
116
117    /// Check if flow is running
118    pub fn is_running(&self) -> bool {
119        matches!(
120            self.status,
121            ExampleSagaFlowStatus::Running | ExampleSagaFlowStatus::Waiting
122        )
123    }
124
125    /// Set a context variable
126    pub fn set_context(&mut self, key: &str, value: serde_json::Value) {
127        if let serde_json::Value::Object(ref mut map) = self.context {
128            map.insert(key.to_string(), value);
129        }
130        self.updated_at = chrono::Utc::now();
131    }
132
133    /// Get a context variable
134    pub fn get_context(&self, key: &str) -> Option<&serde_json::Value> {
135        if let serde_json::Value::Object(ref map) = self.context {
136            map.get(key)
137        } else {
138            None
139        }
140    }
141}
142
143// Phase 1: FlowInstance satisfies backbone_core::flow::WorkflowContext.
144// The entity() accessor requires the entity to be carried in the instance;
145// add `pub entity: ExampleSaga` in the // <<< CUSTOM section and remove the todo!.
146#[allow(unused_variables)]
147impl WorkflowContext<ExampleSagaFlowInstance> for ExampleSagaFlowInstance {
148    fn entity(&self) -> &ExampleSagaFlowInstance { self }
149    fn set_var(&mut self, key: &str, value: serde_json::Value) {
150        self.set_context(key, value);
151    }
152    fn get_var(&self, key: &str) -> Option<&serde_json::Value> {
153        self.get_context(key)
154    }
155}
156
157/// Step handler trait for ExampleSaga flow
158#[async_trait::async_trait]
159pub trait ExampleSagaStepHandler: Send + Sync {
160    /// Handle validate step
161    async fn handle_validate(
162        &self,
163        instance: &mut ExampleSagaFlowInstance,
164    ) -> Result<Option<ExampleSagaFlowStep>, FlowError>;
165
166    /// Handle write_example step
167    async fn handle_write_example(
168        &self,
169        instance: &mut ExampleSagaFlowInstance,
170    ) -> Result<Option<ExampleSagaFlowStep>, FlowError>;
171
172    /// Handle terminal step done
173    async fn handle_done(
174        &self,
175        instance: &mut ExampleSagaFlowInstance,
176    ) -> Result<(), FlowError>;
177
178    /// Handle terminal step rejected
179    async fn handle_rejected(
180        &self,
181        instance: &mut ExampleSagaFlowInstance,
182    ) -> Result<(), FlowError>;
183
184}
185
186/// Executor for ExampleSaga flow
187pub struct ExampleSagaFlowExecutor<H: ExampleSagaStepHandler> {
188    handler: Arc<H>,
189}
190
191impl<H: ExampleSagaStepHandler> ExampleSagaFlowExecutor<H> {
192    /// Create a new flow executor
193    pub fn new(handler: Arc<H>) -> Self {
194        Self { handler }
195    }
196
197    /// Start a new flow instance
198    pub async fn start(&self, instance_id: impl Into<String>) -> Result<ExampleSagaFlowInstance, FlowError> {
199        let mut instance = ExampleSagaFlowInstance::new(instance_id);
200        instance.status = ExampleSagaFlowStatus::Running;
201        instance.current_step = Some(ExampleSagaFlowStep::Validate);
202        Ok(instance)
203    }
204
205    /// Execute the current step
206    pub async fn execute_step(&self, instance: &mut ExampleSagaFlowInstance) -> Result<(), FlowError> {
207        let current_step = match instance.current_step {
208            Some(step) => step,
209            None => return Err(FlowError::NoCurrentStep),
210        };
211
212        let next_step = match current_step {
213            ExampleSagaFlowStep::Validate => {
214                self.handler.handle_validate(instance).await?
215            }
216            ExampleSagaFlowStep::WriteExample => {
217                self.handler.handle_write_example(instance).await?
218            }
219            ExampleSagaFlowStep::Done => {
220                self.handler.handle_done(instance).await?;
221                None // Terminal step
222            }
223            ExampleSagaFlowStep::Rejected => {
224                self.handler.handle_rejected(instance).await?;
225                None // Terminal step
226            }
227        };
228
229        // Mark current step as completed
230        instance.completed_steps.push(current_step);
231
232        // Move to next step
233        match next_step {
234            Some(next) => {
235                instance.current_step = Some(next);
236            }
237            None => {
238                instance.current_step = None;
239                instance.status = ExampleSagaFlowStatus::Completed;
240            }
241        }
242
243        instance.updated_at = chrono::Utc::now();
244        Ok(())
245    }
246
247    /// Run the flow to completion
248    pub async fn run(&self, instance: &mut ExampleSagaFlowInstance) -> Result<(), FlowError> {
249        while !instance.is_complete() && instance.status != ExampleSagaFlowStatus::Waiting {
250            self.execute_step(instance).await?;
251        }
252        Ok(())
253    }
254
255    /// Cancel the flow
256    pub fn cancel(&self, instance: &mut ExampleSagaFlowInstance) {
257        instance.status = ExampleSagaFlowStatus::Cancelled;
258        instance.updated_at = chrono::Utc::now();
259    }
260
261    /// Mark flow as failed
262    pub fn fail(&self, instance: &mut ExampleSagaFlowInstance, error: impl Into<String>) {
263        instance.status = ExampleSagaFlowStatus::Failed;
264        instance.error = Some(error.into());
265        instance.updated_at = chrono::Utc::now();
266    }
267}
268
269#[cfg(test)]
270mod tests {
271    use super::*;
272
273    #[test]
274    fn test_flow_instance_creation() {
275        let instance = ExampleSagaFlowInstance::new("test-1");
276        assert_eq!(instance.id, "test-1");
277        assert_eq!(instance.status, ExampleSagaFlowStatus::Pending);
278        assert!(!instance.is_complete());
279    }
280
281    #[test]
282    fn test_flow_context() {
283        let mut instance = ExampleSagaFlowInstance::new("test-2");
284        instance.set_context("key", serde_json::json!("value"));
285        let value = instance.get_context("key");
286        assert_eq!(value, Some(&serde_json::json!("value")));
287    }
288
289    #[test]
290    fn test_flow_status_transitions() {
291        let mut instance = ExampleSagaFlowInstance::new("test-3");
292        assert!(!instance.is_running());
293
294        instance.status = ExampleSagaFlowStatus::Running;
295        assert!(instance.is_running());
296        assert!(!instance.is_complete());
297
298        instance.status = ExampleSagaFlowStatus::Completed;
299        assert!(instance.is_complete());
300        assert!(!instance.is_running());
301    }
302}