backbone_integrations/application/workflows/
example_saga_workflow.rs1use serde::{Deserialize, Serialize};
9use std::sync::Arc;
10use chrono;
11use backbone_core::flow::{WorkflowStep, WorkflowContext};
12
13#[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#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
40#[serde(rename_all = "snake_case")]
41pub enum ExampleSagaFlowStatus {
42 Pending,
44 Running,
46 Waiting,
48 Completed,
50 Failed,
52 Cancelled,
54 Compensating,
56}
57
58#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
60#[serde(rename_all = "snake_case")]
61pub enum ExampleSagaFlowStep {
62 Validate,
64 WriteExample,
66 Done,
68 Rejected,
70}
71
72#[derive(Debug, Clone, Serialize, Deserialize)]
74pub struct ExampleSagaFlowInstance {
75 pub id: String,
77 pub status: ExampleSagaFlowStatus,
79 pub current_step: Option<ExampleSagaFlowStep>,
81 pub context: serde_json::Value,
83 pub completed_steps: Vec<ExampleSagaFlowStep>,
85 pub error: Option<String>,
87 pub created_at: chrono::DateTime<chrono::Utc>,
89 pub updated_at: chrono::DateTime<chrono::Utc>,
91}
92
93impl ExampleSagaFlowInstance {
94 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 pub fn is_complete(&self) -> bool {
111 matches!(
112 self.status,
113 ExampleSagaFlowStatus::Completed | ExampleSagaFlowStatus::Failed | ExampleSagaFlowStatus::Cancelled
114 )
115 }
116
117 pub fn is_running(&self) -> bool {
119 matches!(
120 self.status,
121 ExampleSagaFlowStatus::Running | ExampleSagaFlowStatus::Waiting
122 )
123 }
124
125 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 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#[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#[async_trait::async_trait]
159pub trait ExampleSagaStepHandler: Send + Sync {
160 async fn handle_validate(
162 &self,
163 instance: &mut ExampleSagaFlowInstance,
164 ) -> Result<Option<ExampleSagaFlowStep>, FlowError>;
165
166 async fn handle_write_example(
168 &self,
169 instance: &mut ExampleSagaFlowInstance,
170 ) -> Result<Option<ExampleSagaFlowStep>, FlowError>;
171
172 async fn handle_done(
174 &self,
175 instance: &mut ExampleSagaFlowInstance,
176 ) -> Result<(), FlowError>;
177
178 async fn handle_rejected(
180 &self,
181 instance: &mut ExampleSagaFlowInstance,
182 ) -> Result<(), FlowError>;
183
184}
185
186pub struct ExampleSagaFlowExecutor<H: ExampleSagaStepHandler> {
188 handler: Arc<H>,
189}
190
191impl<H: ExampleSagaStepHandler> ExampleSagaFlowExecutor<H> {
192 pub fn new(handler: Arc<H>) -> Self {
194 Self { handler }
195 }
196
197 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 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 }
223 ExampleSagaFlowStep::Rejected => {
224 self.handler.handle_rejected(instance).await?;
225 None }
227 };
228
229 instance.completed_steps.push(current_step);
231
232 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 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 pub fn cancel(&self, instance: &mut ExampleSagaFlowInstance) {
257 instance.status = ExampleSagaFlowStatus::Cancelled;
258 instance.updated_at = chrono::Utc::now();
259 }
260
261 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}