backbone-bucket 0.2.1

Bucket Bounded Context: File Storage Module for Backbone Framework
Documentation
//! MediaProcessing flow implementation
//!
//! Generated by metaphor-schema
//!
//! Asynchronous media processing workflow triggered after a file is uploaded.
//! Generates thumbnails and previews for images, videos, and documents.
//! Reusable by any upload flow that emits FileUploadedEvent.

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 MediaProcessing flow
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum MediaProcessingFlowStatus {
    /// 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 MediaProcessing flow
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum MediaProcessingFlowStep {
    /// check_media_type
    CheckMediaType,
    /// create_image_thumbnail
    CreateImageThumbnail,
    /// create_video_thumbnail
    CreateVideoThumbnail,
    /// create_document_preview
    CreateDocumentPreview,
    /// no_processing_needed
    NoProcessingNeeded,
    /// complete
    Complete,
}

/// Instance of MediaProcessing flow execution
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct MediaProcessingFlowInstance {
    /// Unique instance ID
    pub id: String,
    /// Current status
    pub status: MediaProcessingFlowStatus,
    /// Current step
    pub current_step: Option<MediaProcessingFlowStep>,
    /// Flow context (variables)
    pub context: serde_json::Value,
    /// Completed steps
    pub completed_steps: Vec<MediaProcessingFlowStep>,
    /// 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 MediaProcessingFlowInstance {
    /// Create a new flow instance
    pub fn new(id: impl Into<String>) -> Self {
        let now = chrono::Utc::now();
        Self {
            id: id.into(),
            status: MediaProcessingFlowStatus::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,
            MediaProcessingFlowStatus::Completed | MediaProcessingFlowStatus::Failed | MediaProcessingFlowStatus::Cancelled
        )
    }

    /// Check if flow is running
    pub fn is_running(&self) -> bool {
        matches!(
            self.status,
            MediaProcessingFlowStatus::Running | MediaProcessingFlowStatus::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: MediaProcessing` in the // <<< CUSTOM section and remove the todo!.
#[allow(unused_variables)]
impl WorkflowContext<MediaProcessingFlowInstance> for MediaProcessingFlowInstance {
    fn entity(&self) -> &MediaProcessingFlowInstance { 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 MediaProcessing flow
#[async_trait::async_trait]
pub trait MediaProcessingStepHandler: Send + Sync {
    /// Evaluate condition for check_media_type step
    async fn evaluate_check_media_type(
        &self,
        instance: &MediaProcessingFlowInstance,
    ) -> Result<MediaProcessingFlowStep, FlowError>;

    /// Handle create_image_thumbnail step
    async fn handle_create_image_thumbnail(
        &self,
        instance: &mut MediaProcessingFlowInstance,
    ) -> Result<Option<MediaProcessingFlowStep>, FlowError>;

    /// Handle create_video_thumbnail step
    async fn handle_create_video_thumbnail(
        &self,
        instance: &mut MediaProcessingFlowInstance,
    ) -> Result<Option<MediaProcessingFlowStep>, FlowError>;

    /// Handle create_document_preview step
    async fn handle_create_document_preview(
        &self,
        instance: &mut MediaProcessingFlowInstance,
    ) -> Result<Option<MediaProcessingFlowStep>, FlowError>;

    /// Handle terminal step no_processing_needed
    async fn handle_no_processing_needed(
        &self,
        instance: &mut MediaProcessingFlowInstance,
    ) -> Result<(), FlowError>;

    /// Handle terminal step complete
    async fn handle_complete(
        &self,
        instance: &mut MediaProcessingFlowInstance,
    ) -> Result<(), FlowError>;

}

/// Executor for MediaProcessing flow
pub struct MediaProcessingFlowExecutor<H: MediaProcessingStepHandler> {
    handler: Arc<H>,
}

impl<H: MediaProcessingStepHandler> MediaProcessingFlowExecutor<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<MediaProcessingFlowInstance, FlowError> {
        let mut instance = MediaProcessingFlowInstance::new(instance_id);
        instance.status = MediaProcessingFlowStatus::Running;
        instance.current_step = Some(MediaProcessingFlowStep::CheckMediaType);
        Ok(instance)
    }

    /// Execute the current step
    pub async fn execute_step(&self, instance: &mut MediaProcessingFlowInstance) -> Result<(), FlowError> {
        let current_step = match instance.current_step {
            Some(step) => step,
            None => return Err(FlowError::NoCurrentStep),
        };

        let next_step = match current_step {
            MediaProcessingFlowStep::CheckMediaType => {
                let next = self.handler.evaluate_check_media_type(instance).await?;
                Some(next)
            }
            MediaProcessingFlowStep::CreateImageThumbnail => {
                self.handler.handle_create_image_thumbnail(instance).await?
            }
            MediaProcessingFlowStep::CreateVideoThumbnail => {
                self.handler.handle_create_video_thumbnail(instance).await?
            }
            MediaProcessingFlowStep::CreateDocumentPreview => {
                self.handler.handle_create_document_preview(instance).await?
            }
            MediaProcessingFlowStep::NoProcessingNeeded => {
                self.handler.handle_no_processing_needed(instance).await?;
                None // Terminal step
            }
            MediaProcessingFlowStep::Complete => {
                self.handler.handle_complete(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 = MediaProcessingFlowStatus::Completed;
            }
        }

        instance.updated_at = chrono::Utc::now();
        Ok(())
    }

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

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

    /// Mark flow as failed
    pub fn fail(&self, instance: &mut MediaProcessingFlowInstance, error: impl Into<String>) {
        instance.status = MediaProcessingFlowStatus::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 = MediaProcessingFlowInstance::new("test-1");
        assert_eq!(instance.id, "test-1");
        assert_eq!(instance.status, MediaProcessingFlowStatus::Pending);
        assert!(!instance.is_complete());
    }

    #[test]
    fn test_flow_context() {
        let mut instance = MediaProcessingFlowInstance::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 = MediaProcessingFlowInstance::new("test-3");
        assert!(!instance.is_running());

        instance.status = MediaProcessingFlowStatus::Running;
        assert!(instance.is_running());
        assert!(!instance.is_complete());

        instance.status = MediaProcessingFlowStatus::Completed;
        assert!(instance.is_complete());
        assert!(!instance.is_running());
    }
}