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 FileUploadFlowStatus {
Pending,
Running,
Waiting,
Completed,
Failed,
Cancelled,
Compensating,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum FileUploadFlowStep {
LoadBucket,
CheckBucketActive,
CheckDeduplicationEnabled,
CalculateChecksum,
CheckDuplicateContent,
UseExistingContent,
DetectMimeType,
ValidateFileType,
ValidateFileSize,
LoadQuota,
CreateDefaultQuota,
CheckQuota,
CreateFileRecord,
CheckShouldStoreContent,
StoreContent,
CreateContentHash,
LinkContentHash,
UpdateFileProcessing,
ScanFile,
CheckScanResult,
UpdateScanComplete,
CheckCdnEnabled,
GenerateCdnUrl,
UpdateCdnInfo,
MarkFileActive,
UpdateQuota,
UpdateBucketStats,
EmitUploadedEvent,
Complete,
QuarantineFile,
EmitThreatEvent,
NotifyThreat,
CompleteQuarantine,
BucketNotFound,
BucketNotActive,
MimeTypeNotAllowed,
FileTooLarge,
QuotaExceeded,
QuotaExceededTerminal,
RollbackFileRecord,
StorageFailed,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct FileUploadFlowInstance {
pub id: String,
pub status: FileUploadFlowStatus,
pub current_step: Option<FileUploadFlowStep>,
pub context: serde_json::Value,
pub completed_steps: Vec<FileUploadFlowStep>,
pub error: Option<String>,
pub created_at: chrono::DateTime<chrono::Utc>,
pub updated_at: chrono::DateTime<chrono::Utc>,
}
impl FileUploadFlowInstance {
pub fn new(id: impl Into<String>) -> Self {
let now = chrono::Utc::now();
Self {
id: id.into(),
status: FileUploadFlowStatus::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,
FileUploadFlowStatus::Completed | FileUploadFlowStatus::Failed | FileUploadFlowStatus::Cancelled
)
}
pub fn is_running(&self) -> bool {
matches!(
self.status,
FileUploadFlowStatus::Running | FileUploadFlowStatus::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<FileUploadFlowInstance> for FileUploadFlowInstance {
fn entity(&self) -> &FileUploadFlowInstance { 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 FileUploadStepHandler: Send + Sync {
async fn handle_load_bucket(
&self,
instance: &mut FileUploadFlowInstance,
) -> Result<Option<FileUploadFlowStep>, FlowError>;
async fn evaluate_check_bucket_active(
&self,
instance: &FileUploadFlowInstance,
) -> Result<FileUploadFlowStep, FlowError>;
async fn evaluate_check_deduplication_enabled(
&self,
instance: &FileUploadFlowInstance,
) -> Result<FileUploadFlowStep, FlowError>;
async fn handle_calculate_checksum(
&self,
instance: &mut FileUploadFlowInstance,
) -> Result<Option<FileUploadFlowStep>, FlowError>;
async fn handle_check_duplicate_content(
&self,
instance: &mut FileUploadFlowInstance,
) -> Result<Option<FileUploadFlowStep>, FlowError>;
async fn handle_use_existing_content(
&self,
instance: &mut FileUploadFlowInstance,
) -> Result<Option<FileUploadFlowStep>, FlowError>;
async fn handle_detect_mime_type(
&self,
instance: &mut FileUploadFlowInstance,
) -> Result<Option<FileUploadFlowStep>, FlowError>;
async fn evaluate_validate_file_type(
&self,
instance: &FileUploadFlowInstance,
) -> Result<FileUploadFlowStep, FlowError>;
async fn evaluate_validate_file_size(
&self,
instance: &FileUploadFlowInstance,
) -> Result<FileUploadFlowStep, FlowError>;
async fn handle_load_quota(
&self,
instance: &mut FileUploadFlowInstance,
) -> Result<Option<FileUploadFlowStep>, FlowError>;
async fn handle_create_default_quota(
&self,
instance: &mut FileUploadFlowInstance,
) -> Result<Option<FileUploadFlowStep>, FlowError>;
async fn evaluate_check_quota(
&self,
instance: &FileUploadFlowInstance,
) -> Result<FileUploadFlowStep, FlowError>;
async fn handle_create_file_record(
&self,
instance: &mut FileUploadFlowInstance,
) -> Result<Option<FileUploadFlowStep>, FlowError>;
async fn evaluate_check_should_store_content(
&self,
instance: &FileUploadFlowInstance,
) -> Result<FileUploadFlowStep, FlowError>;
async fn handle_store_content(
&self,
instance: &mut FileUploadFlowInstance,
) -> Result<Option<FileUploadFlowStep>, FlowError>;
async fn handle_create_content_hash(
&self,
instance: &mut FileUploadFlowInstance,
) -> Result<Option<FileUploadFlowStep>, FlowError>;
async fn handle_link_content_hash(
&self,
instance: &mut FileUploadFlowInstance,
) -> Result<Option<FileUploadFlowStep>, FlowError>;
async fn handle_update_file_processing(
&self,
instance: &mut FileUploadFlowInstance,
) -> Result<Option<FileUploadFlowStep>, FlowError>;
async fn handle_scan_file(
&self,
instance: &mut FileUploadFlowInstance,
) -> Result<Option<FileUploadFlowStep>, FlowError>;
async fn evaluate_check_scan_result(
&self,
instance: &FileUploadFlowInstance,
) -> Result<FileUploadFlowStep, FlowError>;
async fn handle_update_scan_complete(
&self,
instance: &mut FileUploadFlowInstance,
) -> Result<Option<FileUploadFlowStep>, FlowError>;
async fn evaluate_check_cdn_enabled(
&self,
instance: &FileUploadFlowInstance,
) -> Result<FileUploadFlowStep, FlowError>;
async fn handle_generate_cdn_url(
&self,
instance: &mut FileUploadFlowInstance,
) -> Result<Option<FileUploadFlowStep>, FlowError>;
async fn handle_update_cdn_info(
&self,
instance: &mut FileUploadFlowInstance,
) -> Result<Option<FileUploadFlowStep>, FlowError>;
async fn handle_mark_file_active(
&self,
instance: &mut FileUploadFlowInstance,
) -> Result<Option<FileUploadFlowStep>, FlowError>;
async fn handle_update_quota(
&self,
instance: &mut FileUploadFlowInstance,
) -> Result<Option<FileUploadFlowStep>, FlowError>;
async fn handle_update_bucket_stats(
&self,
instance: &mut FileUploadFlowInstance,
) -> Result<Option<FileUploadFlowStep>, FlowError>;
async fn handle_emit_uploaded_event(
&self,
instance: &mut FileUploadFlowInstance,
) -> Result<Option<FileUploadFlowStep>, FlowError>;
async fn handle_complete(
&self,
instance: &mut FileUploadFlowInstance,
) -> Result<(), FlowError>;
async fn handle_quarantine_file(
&self,
instance: &mut FileUploadFlowInstance,
) -> Result<Option<FileUploadFlowStep>, FlowError>;
async fn handle_emit_threat_event(
&self,
instance: &mut FileUploadFlowInstance,
) -> Result<Option<FileUploadFlowStep>, FlowError>;
async fn handle_notify_threat(
&self,
instance: &mut FileUploadFlowInstance,
) -> Result<Option<FileUploadFlowStep>, FlowError>;
async fn handle_complete_quarantine(
&self,
instance: &mut FileUploadFlowInstance,
) -> Result<(), FlowError>;
async fn handle_bucket_not_found(
&self,
instance: &mut FileUploadFlowInstance,
) -> Result<(), FlowError>;
async fn handle_bucket_not_active(
&self,
instance: &mut FileUploadFlowInstance,
) -> Result<(), FlowError>;
async fn handle_mime_type_not_allowed(
&self,
instance: &mut FileUploadFlowInstance,
) -> Result<(), FlowError>;
async fn handle_file_too_large(
&self,
instance: &mut FileUploadFlowInstance,
) -> Result<(), FlowError>;
async fn handle_quota_exceeded(
&self,
instance: &mut FileUploadFlowInstance,
) -> Result<Option<FileUploadFlowStep>, FlowError>;
async fn handle_quota_exceeded_terminal(
&self,
instance: &mut FileUploadFlowInstance,
) -> Result<(), FlowError>;
async fn handle_rollback_file_record(
&self,
instance: &mut FileUploadFlowInstance,
) -> Result<Option<FileUploadFlowStep>, FlowError>;
async fn handle_storage_failed(
&self,
instance: &mut FileUploadFlowInstance,
) -> Result<(), FlowError>;
async fn compensate(
&self,
instance: &mut FileUploadFlowInstance,
) -> Result<(), FlowError>;
}
pub struct FileUploadFlowExecutor<H: FileUploadStepHandler> {
handler: Arc<H>,
}
impl<H: FileUploadStepHandler> FileUploadFlowExecutor<H> {
pub fn new(handler: Arc<H>) -> Self {
Self { handler }
}
pub async fn start(&self, instance_id: impl Into<String>) -> Result<FileUploadFlowInstance, FlowError> {
let mut instance = FileUploadFlowInstance::new(instance_id);
instance.status = FileUploadFlowStatus::Running;
instance.current_step = Some(FileUploadFlowStep::LoadBucket);
Ok(instance)
}
pub async fn execute_step(&self, instance: &mut FileUploadFlowInstance) -> Result<(), FlowError> {
let current_step = match instance.current_step {
Some(step) => step,
None => return Err(FlowError::NoCurrentStep),
};
let next_step = match current_step {
FileUploadFlowStep::LoadBucket => {
self.handler.handle_load_bucket(instance).await?
}
FileUploadFlowStep::CheckBucketActive => {
let next = self.handler.evaluate_check_bucket_active(instance).await?;
Some(next)
}
FileUploadFlowStep::CheckDeduplicationEnabled => {
let next = self.handler.evaluate_check_deduplication_enabled(instance).await?;
Some(next)
}
FileUploadFlowStep::CalculateChecksum => {
self.handler.handle_calculate_checksum(instance).await?
}
FileUploadFlowStep::CheckDuplicateContent => {
self.handler.handle_check_duplicate_content(instance).await?
}
FileUploadFlowStep::UseExistingContent => {
self.handler.handle_use_existing_content(instance).await?
}
FileUploadFlowStep::DetectMimeType => {
self.handler.handle_detect_mime_type(instance).await?
}
FileUploadFlowStep::ValidateFileType => {
let next = self.handler.evaluate_validate_file_type(instance).await?;
Some(next)
}
FileUploadFlowStep::ValidateFileSize => {
let next = self.handler.evaluate_validate_file_size(instance).await?;
Some(next)
}
FileUploadFlowStep::LoadQuota => {
self.handler.handle_load_quota(instance).await?
}
FileUploadFlowStep::CreateDefaultQuota => {
self.handler.handle_create_default_quota(instance).await?
}
FileUploadFlowStep::CheckQuota => {
let next = self.handler.evaluate_check_quota(instance).await?;
Some(next)
}
FileUploadFlowStep::CreateFileRecord => {
self.handler.handle_create_file_record(instance).await?
}
FileUploadFlowStep::CheckShouldStoreContent => {
let next = self.handler.evaluate_check_should_store_content(instance).await?;
Some(next)
}
FileUploadFlowStep::StoreContent => {
self.handler.handle_store_content(instance).await?
}
FileUploadFlowStep::CreateContentHash => {
self.handler.handle_create_content_hash(instance).await?
}
FileUploadFlowStep::LinkContentHash => {
self.handler.handle_link_content_hash(instance).await?
}
FileUploadFlowStep::UpdateFileProcessing => {
self.handler.handle_update_file_processing(instance).await?
}
FileUploadFlowStep::ScanFile => {
self.handler.handle_scan_file(instance).await?
}
FileUploadFlowStep::CheckScanResult => {
let next = self.handler.evaluate_check_scan_result(instance).await?;
Some(next)
}
FileUploadFlowStep::UpdateScanComplete => {
self.handler.handle_update_scan_complete(instance).await?
}
FileUploadFlowStep::CheckCdnEnabled => {
let next = self.handler.evaluate_check_cdn_enabled(instance).await?;
Some(next)
}
FileUploadFlowStep::GenerateCdnUrl => {
self.handler.handle_generate_cdn_url(instance).await?
}
FileUploadFlowStep::UpdateCdnInfo => {
self.handler.handle_update_cdn_info(instance).await?
}
FileUploadFlowStep::MarkFileActive => {
self.handler.handle_mark_file_active(instance).await?
}
FileUploadFlowStep::UpdateQuota => {
self.handler.handle_update_quota(instance).await?
}
FileUploadFlowStep::UpdateBucketStats => {
self.handler.handle_update_bucket_stats(instance).await?
}
FileUploadFlowStep::EmitUploadedEvent => {
self.handler.handle_emit_uploaded_event(instance).await?
}
FileUploadFlowStep::Complete => {
self.handler.handle_complete(instance).await?;
None }
FileUploadFlowStep::QuarantineFile => {
self.handler.handle_quarantine_file(instance).await?
}
FileUploadFlowStep::EmitThreatEvent => {
self.handler.handle_emit_threat_event(instance).await?
}
FileUploadFlowStep::NotifyThreat => {
self.handler.handle_notify_threat(instance).await?
}
FileUploadFlowStep::CompleteQuarantine => {
self.handler.handle_complete_quarantine(instance).await?;
None }
FileUploadFlowStep::BucketNotFound => {
self.handler.handle_bucket_not_found(instance).await?;
None }
FileUploadFlowStep::BucketNotActive => {
self.handler.handle_bucket_not_active(instance).await?;
None }
FileUploadFlowStep::MimeTypeNotAllowed => {
self.handler.handle_mime_type_not_allowed(instance).await?;
None }
FileUploadFlowStep::FileTooLarge => {
self.handler.handle_file_too_large(instance).await?;
None }
FileUploadFlowStep::QuotaExceeded => {
self.handler.handle_quota_exceeded(instance).await?
}
FileUploadFlowStep::QuotaExceededTerminal => {
self.handler.handle_quota_exceeded_terminal(instance).await?;
None }
FileUploadFlowStep::RollbackFileRecord => {
self.handler.handle_rollback_file_record(instance).await?
}
FileUploadFlowStep::StorageFailed => {
self.handler.handle_storage_failed(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 = FileUploadFlowStatus::Completed;
}
}
instance.updated_at = chrono::Utc::now();
Ok(())
}
pub async fn run(&self, instance: &mut FileUploadFlowInstance) -> Result<(), FlowError> {
while !instance.is_complete() && instance.status != FileUploadFlowStatus::Waiting {
self.execute_step(instance).await?;
}
Ok(())
}
pub fn cancel(&self, instance: &mut FileUploadFlowInstance) {
instance.status = FileUploadFlowStatus::Cancelled;
instance.updated_at = chrono::Utc::now();
}
pub fn fail(&self, instance: &mut FileUploadFlowInstance, error: impl Into<String>) {
instance.status = FileUploadFlowStatus::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 = FileUploadFlowInstance::new("test-1");
assert_eq!(instance.id, "test-1");
assert_eq!(instance.status, FileUploadFlowStatus::Pending);
assert!(!instance.is_complete());
}
#[test]
fn test_flow_context() {
let mut instance = FileUploadFlowInstance::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 = FileUploadFlowInstance::new("test-3");
assert!(!instance.is_running());
instance.status = FileUploadFlowStatus::Running;
assert!(instance.is_running());
assert!(!instance.is_complete());
instance.status = FileUploadFlowStatus::Completed;
assert!(instance.is_complete());
assert!(!instance.is_running());
}
}