backbone-integrations 0.6.0

Integration registry: connectors, integration accounts and an idempotent inbound event lane, with one OAuth flow (HMAC-bound state, PKCE)
Documentation
//! ExampleSaga flow implementation
//!
//! Generated by metaphor-schema
//!
//! Replace with the business intent of your saga — what it creates/updates atomically,
//! what it validates first, and what event it emits on success.

use serde::{Deserialize, Serialize};
use std::sync::Arc;
use chrono;
use backbone_core::flow::{WorkflowStep, WorkflowContext};

/// Error type for flow execution
#[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 },
}

/// Execution status for ExampleSaga flow
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum ExampleSagaFlowStatus {
    /// Flow is pending execution
    Pending,
    /// Flow is currently running
    Running,
    /// Flow is waiting for an event or condition
    Waiting,
    /// Flow completed successfully
    Completed,
    /// Flow failed
    Failed,
    /// Flow was cancelled
    Cancelled,
    /// Flow is compensating (rolling back)
    Compensating,
}

/// Steps in ExampleSaga flow
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum ExampleSagaFlowStep {
    /// validate
    Validate,
    /// write_example
    WriteExample,
    /// done
    Done,
    /// rejected
    Rejected,
}

/// Instance of ExampleSaga flow execution
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ExampleSagaFlowInstance {
    /// Unique instance ID
    pub id: String,
    /// Current status
    pub status: ExampleSagaFlowStatus,
    /// Current step
    pub current_step: Option<ExampleSagaFlowStep>,
    /// Flow context (variables)
    pub context: serde_json::Value,
    /// Completed steps
    pub completed_steps: Vec<ExampleSagaFlowStep>,
    /// Error if failed
    pub error: Option<String>,
    /// Created timestamp
    pub created_at: chrono::DateTime<chrono::Utc>,
    /// Updated timestamp
    pub updated_at: chrono::DateTime<chrono::Utc>,
}

impl ExampleSagaFlowInstance {
    /// Create a new flow instance
    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,
        }
    }

    /// Check if flow is complete
    pub fn is_complete(&self) -> bool {
        matches!(
            self.status,
            ExampleSagaFlowStatus::Completed | ExampleSagaFlowStatus::Failed | ExampleSagaFlowStatus::Cancelled
        )
    }

    /// Check if flow is running
    pub fn is_running(&self) -> bool {
        matches!(
            self.status,
            ExampleSagaFlowStatus::Running | ExampleSagaFlowStatus::Waiting
        )
    }

    /// Set a context variable
    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();
    }

    /// Get a context variable
    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
        }
    }
}

// Phase 1: FlowInstance satisfies backbone_core::flow::WorkflowContext.
// The entity() accessor requires the entity to be carried in the instance;
// add `pub entity: ExampleSaga` in the // <<< CUSTOM section and remove the todo!.
#[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)
    }
}

/// Step handler trait for ExampleSaga flow
#[async_trait::async_trait]
pub trait ExampleSagaStepHandler: Send + Sync {
    /// Handle validate step
    async fn handle_validate(
        &self,
        instance: &mut ExampleSagaFlowInstance,
    ) -> Result<Option<ExampleSagaFlowStep>, FlowError>;

    /// Handle write_example step
    async fn handle_write_example(
        &self,
        instance: &mut ExampleSagaFlowInstance,
    ) -> Result<Option<ExampleSagaFlowStep>, FlowError>;

    /// Handle terminal step done
    async fn handle_done(
        &self,
        instance: &mut ExampleSagaFlowInstance,
    ) -> Result<(), FlowError>;

    /// Handle terminal step rejected
    async fn handle_rejected(
        &self,
        instance: &mut ExampleSagaFlowInstance,
    ) -> Result<(), FlowError>;

}

/// Executor for ExampleSaga flow
pub struct ExampleSagaFlowExecutor<H: ExampleSagaStepHandler> {
    handler: Arc<H>,
}

impl<H: ExampleSagaStepHandler> ExampleSagaFlowExecutor<H> {
    /// Create a new flow executor
    pub fn new(handler: Arc<H>) -> Self {
        Self { handler }
    }

    /// Start a new flow instance
    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)
    }

    /// Execute the current step
    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 // Terminal step
            }
            ExampleSagaFlowStep::Rejected => {
                self.handler.handle_rejected(instance).await?;
                None // Terminal step
            }
        };

        // Mark current step as completed
        instance.completed_steps.push(current_step);

        // Move to next 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(())
    }

    /// Run the flow to completion
    pub async fn run(&self, instance: &mut ExampleSagaFlowInstance) -> Result<(), FlowError> {
        while !instance.is_complete() && instance.status != ExampleSagaFlowStatus::Waiting {
            self.execute_step(instance).await?;
        }
        Ok(())
    }

    /// Cancel the flow
    pub fn cancel(&self, instance: &mut ExampleSagaFlowInstance) {
        instance.status = ExampleSagaFlowStatus::Cancelled;
        instance.updated_at = chrono::Utc::now();
    }

    /// Mark flow as failed
    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());
    }
}