use serde::{Deserialize, Serialize};
use std::sync::Arc;
use chrono;
use backbone_core::flow::{WorkflowStep, WorkflowContext};
#[derive(Debug, Clone, thiserror::Error)]
pub enum FlowError {
#[error("No current step to execute")]
NoCurrentStep,
#[error("Step execution failed: {0}")]
StepFailed(String),
#[error("Condition evaluation failed: {0}")]
ConditionFailed(String),
#[error("Compensation failed: {0}")]
CompensationFailed(String),
#[error("Flow timed out")]
Timeout,
#[error("Flow cancelled")]
Cancelled,
#[error("Invalid state transition: {from} -> {to}")]
InvalidTransition { from: String, to: String },
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum ExampleSagaFlowStatus {
Pending,
Running,
Waiting,
Completed,
Failed,
Cancelled,
Compensating,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum ExampleSagaFlowStep {
Validate,
WriteExample,
Done,
Rejected,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ExampleSagaFlowInstance {
pub id: String,
pub status: ExampleSagaFlowStatus,
pub current_step: Option<ExampleSagaFlowStep>,
pub context: serde_json::Value,
pub completed_steps: Vec<ExampleSagaFlowStep>,
pub error: Option<String>,
pub created_at: chrono::DateTime<chrono::Utc>,
pub updated_at: chrono::DateTime<chrono::Utc>,
}
impl ExampleSagaFlowInstance {
pub fn new(id: impl Into<String>) -> Self {
let now = chrono::Utc::now();
Self {
id: id.into(),
status: ExampleSagaFlowStatus::Pending,
current_step: None,
context: serde_json::json!({}),
completed_steps: Vec::new(),
error: None,
created_at: now,
updated_at: now,
}
}
pub fn is_complete(&self) -> bool {
matches!(
self.status,
ExampleSagaFlowStatus::Completed | ExampleSagaFlowStatus::Failed | ExampleSagaFlowStatus::Cancelled
)
}
pub fn is_running(&self) -> bool {
matches!(
self.status,
ExampleSagaFlowStatus::Running | ExampleSagaFlowStatus::Waiting
)
}
pub fn set_context(&mut self, key: &str, value: serde_json::Value) {
if let serde_json::Value::Object(ref mut map) = self.context {
map.insert(key.to_string(), value);
}
self.updated_at = chrono::Utc::now();
}
pub fn get_context(&self, key: &str) -> Option<&serde_json::Value> {
if let serde_json::Value::Object(ref map) = self.context {
map.get(key)
} else {
None
}
}
}
#[allow(unused_variables)]
impl WorkflowContext<ExampleSagaFlowInstance> for ExampleSagaFlowInstance {
fn entity(&self) -> &ExampleSagaFlowInstance { self }
fn set_var(&mut self, key: &str, value: serde_json::Value) {
self.set_context(key, value);
}
fn get_var(&self, key: &str) -> Option<&serde_json::Value> {
self.get_context(key)
}
}
#[async_trait::async_trait]
pub trait ExampleSagaStepHandler: Send + Sync {
async fn handle_validate(
&self,
instance: &mut ExampleSagaFlowInstance,
) -> Result<Option<ExampleSagaFlowStep>, FlowError>;
async fn handle_write_example(
&self,
instance: &mut ExampleSagaFlowInstance,
) -> Result<Option<ExampleSagaFlowStep>, FlowError>;
async fn handle_done(
&self,
instance: &mut ExampleSagaFlowInstance,
) -> Result<(), FlowError>;
async fn handle_rejected(
&self,
instance: &mut ExampleSagaFlowInstance,
) -> Result<(), FlowError>;
}
pub struct ExampleSagaFlowExecutor<H: ExampleSagaStepHandler> {
handler: Arc<H>,
}
impl<H: ExampleSagaStepHandler> ExampleSagaFlowExecutor<H> {
pub fn new(handler: Arc<H>) -> Self {
Self { handler }
}
pub async fn start(&self, instance_id: impl Into<String>) -> Result<ExampleSagaFlowInstance, FlowError> {
let mut instance = ExampleSagaFlowInstance::new(instance_id);
instance.status = ExampleSagaFlowStatus::Running;
instance.current_step = Some(ExampleSagaFlowStep::Validate);
Ok(instance)
}
pub async fn execute_step(&self, instance: &mut ExampleSagaFlowInstance) -> Result<(), FlowError> {
let current_step = match instance.current_step {
Some(step) => step,
None => return Err(FlowError::NoCurrentStep),
};
let next_step = match current_step {
ExampleSagaFlowStep::Validate => {
self.handler.handle_validate(instance).await?
}
ExampleSagaFlowStep::WriteExample => {
self.handler.handle_write_example(instance).await?
}
ExampleSagaFlowStep::Done => {
self.handler.handle_done(instance).await?;
None }
ExampleSagaFlowStep::Rejected => {
self.handler.handle_rejected(instance).await?;
None }
};
instance.completed_steps.push(current_step);
match next_step {
Some(next) => {
instance.current_step = Some(next);
}
None => {
instance.current_step = None;
instance.status = ExampleSagaFlowStatus::Completed;
}
}
instance.updated_at = chrono::Utc::now();
Ok(())
}
pub async fn run(&self, instance: &mut ExampleSagaFlowInstance) -> Result<(), FlowError> {
while !instance.is_complete() && instance.status != ExampleSagaFlowStatus::Waiting {
self.execute_step(instance).await?;
}
Ok(())
}
pub fn cancel(&self, instance: &mut ExampleSagaFlowInstance) {
instance.status = ExampleSagaFlowStatus::Cancelled;
instance.updated_at = chrono::Utc::now();
}
pub fn fail(&self, instance: &mut ExampleSagaFlowInstance, error: impl Into<String>) {
instance.status = ExampleSagaFlowStatus::Failed;
instance.error = Some(error.into());
instance.updated_at = chrono::Utc::now();
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_flow_instance_creation() {
let instance = ExampleSagaFlowInstance::new("test-1");
assert_eq!(instance.id, "test-1");
assert_eq!(instance.status, ExampleSagaFlowStatus::Pending);
assert!(!instance.is_complete());
}
#[test]
fn test_flow_context() {
let mut instance = ExampleSagaFlowInstance::new("test-2");
instance.set_context("key", serde_json::json!("value"));
let value = instance.get_context("key");
assert_eq!(value, Some(&serde_json::json!("value")));
}
#[test]
fn test_flow_status_transitions() {
let mut instance = ExampleSagaFlowInstance::new("test-3");
assert!(!instance.is_running());
instance.status = ExampleSagaFlowStatus::Running;
assert!(instance.is_running());
assert!(!instance.is_complete());
instance.status = ExampleSagaFlowStatus::Completed;
assert!(instance.is_complete());
assert!(!instance.is_running());
}
}