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 FileProcessingFlowStatus {
Pending,
Running,
Waiting,
Completed,
Failed,
Cancelled,
Compensating,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum FileProcessingFlowStep {
LoadJob,
MarkRunning,
LoadFile,
LoadBucket,
RouteByJobType,
ProcessImageThumbnail,
StoreImageThumbnail,
UpdateFileThumbnail,
ProcessVideoThumbnail,
StoreVideoThumbnail,
UpdateFileVideoThumbnail,
ProcessDocumentPreview,
StoreDocumentPreview,
UpdateFileDocumentPreview,
ProcessCompression,
StoreCompressedFile,
UpdateFileCompression,
UpdateBucketCompressionStats,
ProcessVirusScan,
CheckThreatLevel,
UpdateScanSafe,
UpdateScanWarning,
NotifyScanWarning,
UpdateScanThreat,
NotifyThreatDetected,
ProcessDeduplication,
CheckExistingHash,
UpdateExistingHash,
LinkFileHash,
CreateNewHash,
MarkComplete,
EmitCompleteEvent,
Complete,
RetryOrFail,
ScheduleRetry,
RetryTerminal,
MarkFailed,
EmitFailedEvent,
FailedTerminal,
JobNotFound,
FileNotFound,
UnknownJobType,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct FileProcessingFlowInstance {
pub id: String,
pub status: FileProcessingFlowStatus,
pub current_step: Option<FileProcessingFlowStep>,
pub context: serde_json::Value,
pub completed_steps: Vec<FileProcessingFlowStep>,
pub error: Option<String>,
pub created_at: chrono::DateTime<chrono::Utc>,
pub updated_at: chrono::DateTime<chrono::Utc>,
}
impl FileProcessingFlowInstance {
pub fn new(id: impl Into<String>) -> Self {
let now = chrono::Utc::now();
Self {
id: id.into(),
status: FileProcessingFlowStatus::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,
FileProcessingFlowStatus::Completed | FileProcessingFlowStatus::Failed | FileProcessingFlowStatus::Cancelled
)
}
pub fn is_running(&self) -> bool {
matches!(
self.status,
FileProcessingFlowStatus::Running | FileProcessingFlowStatus::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<FileProcessingFlowInstance> for FileProcessingFlowInstance {
fn entity(&self) -> &FileProcessingFlowInstance { 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 FileProcessingStepHandler: Send + Sync {
async fn handle_load_job(
&self,
instance: &mut FileProcessingFlowInstance,
) -> Result<Option<FileProcessingFlowStep>, FlowError>;
async fn handle_mark_running(
&self,
instance: &mut FileProcessingFlowInstance,
) -> Result<Option<FileProcessingFlowStep>, FlowError>;
async fn handle_load_file(
&self,
instance: &mut FileProcessingFlowInstance,
) -> Result<Option<FileProcessingFlowStep>, FlowError>;
async fn handle_load_bucket(
&self,
instance: &mut FileProcessingFlowInstance,
) -> Result<Option<FileProcessingFlowStep>, FlowError>;
async fn evaluate_route_by_job_type(
&self,
instance: &FileProcessingFlowInstance,
) -> Result<FileProcessingFlowStep, FlowError>;
async fn handle_process_image_thumbnail(
&self,
instance: &mut FileProcessingFlowInstance,
) -> Result<Option<FileProcessingFlowStep>, FlowError>;
async fn handle_store_image_thumbnail(
&self,
instance: &mut FileProcessingFlowInstance,
) -> Result<Option<FileProcessingFlowStep>, FlowError>;
async fn handle_update_file_thumbnail(
&self,
instance: &mut FileProcessingFlowInstance,
) -> Result<Option<FileProcessingFlowStep>, FlowError>;
async fn handle_process_video_thumbnail(
&self,
instance: &mut FileProcessingFlowInstance,
) -> Result<Option<FileProcessingFlowStep>, FlowError>;
async fn handle_store_video_thumbnail(
&self,
instance: &mut FileProcessingFlowInstance,
) -> Result<Option<FileProcessingFlowStep>, FlowError>;
async fn handle_update_file_video_thumbnail(
&self,
instance: &mut FileProcessingFlowInstance,
) -> Result<Option<FileProcessingFlowStep>, FlowError>;
async fn handle_process_document_preview(
&self,
instance: &mut FileProcessingFlowInstance,
) -> Result<Option<FileProcessingFlowStep>, FlowError>;
async fn handle_store_document_preview(
&self,
instance: &mut FileProcessingFlowInstance,
) -> Result<Option<FileProcessingFlowStep>, FlowError>;
async fn handle_update_file_document_preview(
&self,
instance: &mut FileProcessingFlowInstance,
) -> Result<Option<FileProcessingFlowStep>, FlowError>;
async fn handle_process_compression(
&self,
instance: &mut FileProcessingFlowInstance,
) -> Result<Option<FileProcessingFlowStep>, FlowError>;
async fn handle_store_compressed_file(
&self,
instance: &mut FileProcessingFlowInstance,
) -> Result<Option<FileProcessingFlowStep>, FlowError>;
async fn handle_update_file_compression(
&self,
instance: &mut FileProcessingFlowInstance,
) -> Result<Option<FileProcessingFlowStep>, FlowError>;
async fn handle_update_bucket_compression_stats(
&self,
instance: &mut FileProcessingFlowInstance,
) -> Result<Option<FileProcessingFlowStep>, FlowError>;
async fn handle_process_virus_scan(
&self,
instance: &mut FileProcessingFlowInstance,
) -> Result<Option<FileProcessingFlowStep>, FlowError>;
async fn evaluate_check_threat_level(
&self,
instance: &FileProcessingFlowInstance,
) -> Result<FileProcessingFlowStep, FlowError>;
async fn handle_update_scan_safe(
&self,
instance: &mut FileProcessingFlowInstance,
) -> Result<Option<FileProcessingFlowStep>, FlowError>;
async fn handle_update_scan_warning(
&self,
instance: &mut FileProcessingFlowInstance,
) -> Result<Option<FileProcessingFlowStep>, FlowError>;
async fn handle_notify_scan_warning(
&self,
instance: &mut FileProcessingFlowInstance,
) -> Result<Option<FileProcessingFlowStep>, FlowError>;
async fn handle_update_scan_threat(
&self,
instance: &mut FileProcessingFlowInstance,
) -> Result<Option<FileProcessingFlowStep>, FlowError>;
async fn handle_notify_threat_detected(
&self,
instance: &mut FileProcessingFlowInstance,
) -> Result<Option<FileProcessingFlowStep>, FlowError>;
async fn handle_process_deduplication(
&self,
instance: &mut FileProcessingFlowInstance,
) -> Result<Option<FileProcessingFlowStep>, FlowError>;
async fn handle_check_existing_hash(
&self,
instance: &mut FileProcessingFlowInstance,
) -> Result<Option<FileProcessingFlowStep>, FlowError>;
async fn handle_update_existing_hash(
&self,
instance: &mut FileProcessingFlowInstance,
) -> Result<Option<FileProcessingFlowStep>, FlowError>;
async fn handle_link_file_hash(
&self,
instance: &mut FileProcessingFlowInstance,
) -> Result<Option<FileProcessingFlowStep>, FlowError>;
async fn handle_create_new_hash(
&self,
instance: &mut FileProcessingFlowInstance,
) -> Result<Option<FileProcessingFlowStep>, FlowError>;
async fn handle_mark_complete(
&self,
instance: &mut FileProcessingFlowInstance,
) -> Result<Option<FileProcessingFlowStep>, FlowError>;
async fn handle_emit_complete_event(
&self,
instance: &mut FileProcessingFlowInstance,
) -> Result<Option<FileProcessingFlowStep>, FlowError>;
async fn handle_complete(
&self,
instance: &mut FileProcessingFlowInstance,
) -> Result<(), FlowError>;
async fn evaluate_retry_or_fail(
&self,
instance: &FileProcessingFlowInstance,
) -> Result<FileProcessingFlowStep, FlowError>;
async fn handle_schedule_retry(
&self,
instance: &mut FileProcessingFlowInstance,
) -> Result<Option<FileProcessingFlowStep>, FlowError>;
async fn handle_retry_terminal(
&self,
instance: &mut FileProcessingFlowInstance,
) -> Result<(), FlowError>;
async fn handle_mark_failed(
&self,
instance: &mut FileProcessingFlowInstance,
) -> Result<Option<FileProcessingFlowStep>, FlowError>;
async fn handle_emit_failed_event(
&self,
instance: &mut FileProcessingFlowInstance,
) -> Result<Option<FileProcessingFlowStep>, FlowError>;
async fn handle_failed_terminal(
&self,
instance: &mut FileProcessingFlowInstance,
) -> Result<(), FlowError>;
async fn handle_job_not_found(
&self,
instance: &mut FileProcessingFlowInstance,
) -> Result<(), FlowError>;
async fn handle_file_not_found(
&self,
instance: &mut FileProcessingFlowInstance,
) -> Result<(), FlowError>;
async fn handle_unknown_job_type(
&self,
instance: &mut FileProcessingFlowInstance,
) -> Result<(), FlowError>;
async fn compensate(
&self,
instance: &mut FileProcessingFlowInstance,
) -> Result<(), FlowError>;
}
pub struct FileProcessingFlowExecutor<H: FileProcessingStepHandler> {
handler: Arc<H>,
}
impl<H: FileProcessingStepHandler> FileProcessingFlowExecutor<H> {
pub fn new(handler: Arc<H>) -> Self {
Self { handler }
}
pub async fn start(&self, instance_id: impl Into<String>) -> Result<FileProcessingFlowInstance, FlowError> {
let mut instance = FileProcessingFlowInstance::new(instance_id);
instance.status = FileProcessingFlowStatus::Running;
instance.current_step = Some(FileProcessingFlowStep::LoadJob);
Ok(instance)
}
pub async fn execute_step(&self, instance: &mut FileProcessingFlowInstance) -> Result<(), FlowError> {
let current_step = match instance.current_step {
Some(step) => step,
None => return Err(FlowError::NoCurrentStep),
};
let next_step = match current_step {
FileProcessingFlowStep::LoadJob => {
self.handler.handle_load_job(instance).await?
}
FileProcessingFlowStep::MarkRunning => {
self.handler.handle_mark_running(instance).await?
}
FileProcessingFlowStep::LoadFile => {
self.handler.handle_load_file(instance).await?
}
FileProcessingFlowStep::LoadBucket => {
self.handler.handle_load_bucket(instance).await?
}
FileProcessingFlowStep::RouteByJobType => {
let next = self.handler.evaluate_route_by_job_type(instance).await?;
Some(next)
}
FileProcessingFlowStep::ProcessImageThumbnail => {
self.handler.handle_process_image_thumbnail(instance).await?
}
FileProcessingFlowStep::StoreImageThumbnail => {
self.handler.handle_store_image_thumbnail(instance).await?
}
FileProcessingFlowStep::UpdateFileThumbnail => {
self.handler.handle_update_file_thumbnail(instance).await?
}
FileProcessingFlowStep::ProcessVideoThumbnail => {
self.handler.handle_process_video_thumbnail(instance).await?
}
FileProcessingFlowStep::StoreVideoThumbnail => {
self.handler.handle_store_video_thumbnail(instance).await?
}
FileProcessingFlowStep::UpdateFileVideoThumbnail => {
self.handler.handle_update_file_video_thumbnail(instance).await?
}
FileProcessingFlowStep::ProcessDocumentPreview => {
self.handler.handle_process_document_preview(instance).await?
}
FileProcessingFlowStep::StoreDocumentPreview => {
self.handler.handle_store_document_preview(instance).await?
}
FileProcessingFlowStep::UpdateFileDocumentPreview => {
self.handler.handle_update_file_document_preview(instance).await?
}
FileProcessingFlowStep::ProcessCompression => {
self.handler.handle_process_compression(instance).await?
}
FileProcessingFlowStep::StoreCompressedFile => {
self.handler.handle_store_compressed_file(instance).await?
}
FileProcessingFlowStep::UpdateFileCompression => {
self.handler.handle_update_file_compression(instance).await?
}
FileProcessingFlowStep::UpdateBucketCompressionStats => {
self.handler.handle_update_bucket_compression_stats(instance).await?
}
FileProcessingFlowStep::ProcessVirusScan => {
self.handler.handle_process_virus_scan(instance).await?
}
FileProcessingFlowStep::CheckThreatLevel => {
let next = self.handler.evaluate_check_threat_level(instance).await?;
Some(next)
}
FileProcessingFlowStep::UpdateScanSafe => {
self.handler.handle_update_scan_safe(instance).await?
}
FileProcessingFlowStep::UpdateScanWarning => {
self.handler.handle_update_scan_warning(instance).await?
}
FileProcessingFlowStep::NotifyScanWarning => {
self.handler.handle_notify_scan_warning(instance).await?
}
FileProcessingFlowStep::UpdateScanThreat => {
self.handler.handle_update_scan_threat(instance).await?
}
FileProcessingFlowStep::NotifyThreatDetected => {
self.handler.handle_notify_threat_detected(instance).await?
}
FileProcessingFlowStep::ProcessDeduplication => {
self.handler.handle_process_deduplication(instance).await?
}
FileProcessingFlowStep::CheckExistingHash => {
self.handler.handle_check_existing_hash(instance).await?
}
FileProcessingFlowStep::UpdateExistingHash => {
self.handler.handle_update_existing_hash(instance).await?
}
FileProcessingFlowStep::LinkFileHash => {
self.handler.handle_link_file_hash(instance).await?
}
FileProcessingFlowStep::CreateNewHash => {
self.handler.handle_create_new_hash(instance).await?
}
FileProcessingFlowStep::MarkComplete => {
self.handler.handle_mark_complete(instance).await?
}
FileProcessingFlowStep::EmitCompleteEvent => {
self.handler.handle_emit_complete_event(instance).await?
}
FileProcessingFlowStep::Complete => {
self.handler.handle_complete(instance).await?;
None }
FileProcessingFlowStep::RetryOrFail => {
let next = self.handler.evaluate_retry_or_fail(instance).await?;
Some(next)
}
FileProcessingFlowStep::ScheduleRetry => {
self.handler.handle_schedule_retry(instance).await?
}
FileProcessingFlowStep::RetryTerminal => {
self.handler.handle_retry_terminal(instance).await?;
None }
FileProcessingFlowStep::MarkFailed => {
self.handler.handle_mark_failed(instance).await?
}
FileProcessingFlowStep::EmitFailedEvent => {
self.handler.handle_emit_failed_event(instance).await?
}
FileProcessingFlowStep::FailedTerminal => {
self.handler.handle_failed_terminal(instance).await?;
None }
FileProcessingFlowStep::JobNotFound => {
self.handler.handle_job_not_found(instance).await?;
None }
FileProcessingFlowStep::FileNotFound => {
self.handler.handle_file_not_found(instance).await?;
None }
FileProcessingFlowStep::UnknownJobType => {
self.handler.handle_unknown_job_type(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 = FileProcessingFlowStatus::Completed;
}
}
instance.updated_at = chrono::Utc::now();
Ok(())
}
pub async fn run(&self, instance: &mut FileProcessingFlowInstance) -> Result<(), FlowError> {
while !instance.is_complete() && instance.status != FileProcessingFlowStatus::Waiting {
self.execute_step(instance).await?;
}
Ok(())
}
pub fn cancel(&self, instance: &mut FileProcessingFlowInstance) {
instance.status = FileProcessingFlowStatus::Cancelled;
instance.updated_at = chrono::Utc::now();
}
pub fn fail(&self, instance: &mut FileProcessingFlowInstance, error: impl Into<String>) {
instance.status = FileProcessingFlowStatus::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 = FileProcessingFlowInstance::new("test-1");
assert_eq!(instance.id, "test-1");
assert_eq!(instance.status, FileProcessingFlowStatus::Pending);
assert!(!instance.is_complete());
}
#[test]
fn test_flow_context() {
let mut instance = FileProcessingFlowInstance::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 = FileProcessingFlowInstance::new("test-3");
assert!(!instance.is_running());
instance.status = FileProcessingFlowStatus::Running;
assert!(instance.is_running());
assert!(!instance.is_complete());
instance.status = FileProcessingFlowStatus::Completed;
assert!(instance.is_complete());
assert!(!instance.is_running());
}
}