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 MultipartUploadFlowStatus {
Pending,
Running,
Waiting,
Completed,
Failed,
Cancelled,
Compensating,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum MultipartUploadFlowStep {
LoadBucket,
CheckBucketAccessible,
ValidateFileSize,
CheckMaxFileSize,
CalculateChunkParams,
ValidateMimeType,
LoadQuota,
CreateDefaultQuota,
CheckQuotaAvailable,
CreateUploadSession,
ReserveQuota,
EmitSessionCreated,
Complete,
BucketNotFound,
BucketNotActive,
InvalidFileSize,
FileTooLarge,
MimeTypeNotAllowed,
QuotaExceeded,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct MultipartUploadFlowInstance {
pub id: String,
pub status: MultipartUploadFlowStatus,
pub current_step: Option<MultipartUploadFlowStep>,
pub context: serde_json::Value,
pub completed_steps: Vec<MultipartUploadFlowStep>,
pub error: Option<String>,
pub created_at: chrono::DateTime<chrono::Utc>,
pub updated_at: chrono::DateTime<chrono::Utc>,
}
impl MultipartUploadFlowInstance {
pub fn new(id: impl Into<String>) -> Self {
let now = chrono::Utc::now();
Self {
id: id.into(),
status: MultipartUploadFlowStatus::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,
MultipartUploadFlowStatus::Completed | MultipartUploadFlowStatus::Failed | MultipartUploadFlowStatus::Cancelled
)
}
pub fn is_running(&self) -> bool {
matches!(
self.status,
MultipartUploadFlowStatus::Running | MultipartUploadFlowStatus::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<MultipartUploadFlowInstance> for MultipartUploadFlowInstance {
fn entity(&self) -> &MultipartUploadFlowInstance { 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 MultipartUploadStepHandler: Send + Sync {
async fn handle_load_bucket(
&self,
instance: &mut MultipartUploadFlowInstance,
) -> Result<Option<MultipartUploadFlowStep>, FlowError>;
async fn evaluate_check_bucket_accessible(
&self,
instance: &MultipartUploadFlowInstance,
) -> Result<MultipartUploadFlowStep, FlowError>;
async fn evaluate_validate_file_size(
&self,
instance: &MultipartUploadFlowInstance,
) -> Result<MultipartUploadFlowStep, FlowError>;
async fn evaluate_check_max_file_size(
&self,
instance: &MultipartUploadFlowInstance,
) -> Result<MultipartUploadFlowStep, FlowError>;
async fn handle_calculate_chunk_params(
&self,
instance: &mut MultipartUploadFlowInstance,
) -> Result<Option<MultipartUploadFlowStep>, FlowError>;
async fn evaluate_validate_mime_type(
&self,
instance: &MultipartUploadFlowInstance,
) -> Result<MultipartUploadFlowStep, FlowError>;
async fn handle_load_quota(
&self,
instance: &mut MultipartUploadFlowInstance,
) -> Result<Option<MultipartUploadFlowStep>, FlowError>;
async fn handle_create_default_quota(
&self,
instance: &mut MultipartUploadFlowInstance,
) -> Result<Option<MultipartUploadFlowStep>, FlowError>;
async fn evaluate_check_quota_available(
&self,
instance: &MultipartUploadFlowInstance,
) -> Result<MultipartUploadFlowStep, FlowError>;
async fn handle_create_upload_session(
&self,
instance: &mut MultipartUploadFlowInstance,
) -> Result<Option<MultipartUploadFlowStep>, FlowError>;
async fn handle_reserve_quota(
&self,
instance: &mut MultipartUploadFlowInstance,
) -> Result<Option<MultipartUploadFlowStep>, FlowError>;
async fn handle_emit_session_created(
&self,
instance: &mut MultipartUploadFlowInstance,
) -> Result<Option<MultipartUploadFlowStep>, FlowError>;
async fn handle_complete(
&self,
instance: &mut MultipartUploadFlowInstance,
) -> Result<(), FlowError>;
async fn handle_bucket_not_found(
&self,
instance: &mut MultipartUploadFlowInstance,
) -> Result<(), FlowError>;
async fn handle_bucket_not_active(
&self,
instance: &mut MultipartUploadFlowInstance,
) -> Result<(), FlowError>;
async fn handle_invalid_file_size(
&self,
instance: &mut MultipartUploadFlowInstance,
) -> Result<(), FlowError>;
async fn handle_file_too_large(
&self,
instance: &mut MultipartUploadFlowInstance,
) -> Result<(), FlowError>;
async fn handle_mime_type_not_allowed(
&self,
instance: &mut MultipartUploadFlowInstance,
) -> Result<(), FlowError>;
async fn handle_quota_exceeded(
&self,
instance: &mut MultipartUploadFlowInstance,
) -> Result<(), FlowError>;
}
pub struct MultipartUploadFlowExecutor<H: MultipartUploadStepHandler> {
handler: Arc<H>,
}
impl<H: MultipartUploadStepHandler> MultipartUploadFlowExecutor<H> {
pub fn new(handler: Arc<H>) -> Self {
Self { handler }
}
pub async fn start(&self, instance_id: impl Into<String>) -> Result<MultipartUploadFlowInstance, FlowError> {
let mut instance = MultipartUploadFlowInstance::new(instance_id);
instance.status = MultipartUploadFlowStatus::Running;
instance.current_step = Some(MultipartUploadFlowStep::LoadBucket);
Ok(instance)
}
pub async fn execute_step(&self, instance: &mut MultipartUploadFlowInstance) -> Result<(), FlowError> {
let current_step = match instance.current_step {
Some(step) => step,
None => return Err(FlowError::NoCurrentStep),
};
let next_step = match current_step {
MultipartUploadFlowStep::LoadBucket => {
self.handler.handle_load_bucket(instance).await?
}
MultipartUploadFlowStep::CheckBucketAccessible => {
let next = self.handler.evaluate_check_bucket_accessible(instance).await?;
Some(next)
}
MultipartUploadFlowStep::ValidateFileSize => {
let next = self.handler.evaluate_validate_file_size(instance).await?;
Some(next)
}
MultipartUploadFlowStep::CheckMaxFileSize => {
let next = self.handler.evaluate_check_max_file_size(instance).await?;
Some(next)
}
MultipartUploadFlowStep::CalculateChunkParams => {
self.handler.handle_calculate_chunk_params(instance).await?
}
MultipartUploadFlowStep::ValidateMimeType => {
let next = self.handler.evaluate_validate_mime_type(instance).await?;
Some(next)
}
MultipartUploadFlowStep::LoadQuota => {
self.handler.handle_load_quota(instance).await?
}
MultipartUploadFlowStep::CreateDefaultQuota => {
self.handler.handle_create_default_quota(instance).await?
}
MultipartUploadFlowStep::CheckQuotaAvailable => {
let next = self.handler.evaluate_check_quota_available(instance).await?;
Some(next)
}
MultipartUploadFlowStep::CreateUploadSession => {
self.handler.handle_create_upload_session(instance).await?
}
MultipartUploadFlowStep::ReserveQuota => {
self.handler.handle_reserve_quota(instance).await?
}
MultipartUploadFlowStep::EmitSessionCreated => {
self.handler.handle_emit_session_created(instance).await?
}
MultipartUploadFlowStep::Complete => {
self.handler.handle_complete(instance).await?;
None }
MultipartUploadFlowStep::BucketNotFound => {
self.handler.handle_bucket_not_found(instance).await?;
None }
MultipartUploadFlowStep::BucketNotActive => {
self.handler.handle_bucket_not_active(instance).await?;
None }
MultipartUploadFlowStep::InvalidFileSize => {
self.handler.handle_invalid_file_size(instance).await?;
None }
MultipartUploadFlowStep::FileTooLarge => {
self.handler.handle_file_too_large(instance).await?;
None }
MultipartUploadFlowStep::MimeTypeNotAllowed => {
self.handler.handle_mime_type_not_allowed(instance).await?;
None }
MultipartUploadFlowStep::QuotaExceeded => {
self.handler.handle_quota_exceeded(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 = MultipartUploadFlowStatus::Completed;
}
}
instance.updated_at = chrono::Utc::now();
Ok(())
}
pub async fn run(&self, instance: &mut MultipartUploadFlowInstance) -> Result<(), FlowError> {
while !instance.is_complete() && instance.status != MultipartUploadFlowStatus::Waiting {
self.execute_step(instance).await?;
}
Ok(())
}
pub fn cancel(&self, instance: &mut MultipartUploadFlowInstance) {
instance.status = MultipartUploadFlowStatus::Cancelled;
instance.updated_at = chrono::Utc::now();
}
pub fn fail(&self, instance: &mut MultipartUploadFlowInstance, error: impl Into<String>) {
instance.status = MultipartUploadFlowStatus::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 = MultipartUploadFlowInstance::new("test-1");
assert_eq!(instance.id, "test-1");
assert_eq!(instance.status, MultipartUploadFlowStatus::Pending);
assert!(!instance.is_complete());
}
#[test]
fn test_flow_context() {
let mut instance = MultipartUploadFlowInstance::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 = MultipartUploadFlowInstance::new("test-3");
assert!(!instance.is_running());
instance.status = MultipartUploadFlowStatus::Running;
assert!(instance.is_running());
assert!(!instance.is_complete());
instance.status = MultipartUploadFlowStatus::Completed;
assert!(instance.is_complete());
assert!(!instance.is_running());
}
}