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 ShareCreationFlowStatus {
Pending,
Running,
Waiting,
Completed,
Failed,
Cancelled,
Compensating,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum ShareCreationFlowStep {
ValidateFile,
CheckSharePermission,
CheckManagePermission,
CheckShareType,
ValidatePublicShare,
ValidateUserShare,
ValidateEmailShare,
CheckExistingShare,
UpdateExistingShare,
CreateShareRecord,
GenerateShareUrl,
DetermineNotification,
NotifyUser,
SendEmailInvitation,
LogShare,
Complete,
FileNotFound,
PermissionDenied,
FailedPermission,
PublicSharingDisabled,
UserNotFound,
InvalidConfig,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ShareCreationFlowInstance {
pub id: String,
pub status: ShareCreationFlowStatus,
pub current_step: Option<ShareCreationFlowStep>,
pub context: serde_json::Value,
pub completed_steps: Vec<ShareCreationFlowStep>,
pub error: Option<String>,
pub created_at: chrono::DateTime<chrono::Utc>,
pub updated_at: chrono::DateTime<chrono::Utc>,
}
impl ShareCreationFlowInstance {
pub fn new(id: impl Into<String>) -> Self {
let now = chrono::Utc::now();
Self {
id: id.into(),
status: ShareCreationFlowStatus::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,
ShareCreationFlowStatus::Completed | ShareCreationFlowStatus::Failed | ShareCreationFlowStatus::Cancelled
)
}
pub fn is_running(&self) -> bool {
matches!(
self.status,
ShareCreationFlowStatus::Running | ShareCreationFlowStatus::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<ShareCreationFlowInstance> for ShareCreationFlowInstance {
fn entity(&self) -> &ShareCreationFlowInstance { 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 ShareCreationStepHandler: Send + Sync {
async fn handle_validate_file(
&self,
instance: &mut ShareCreationFlowInstance,
) -> Result<Option<ShareCreationFlowStep>, FlowError>;
async fn evaluate_check_share_permission(
&self,
instance: &ShareCreationFlowInstance,
) -> Result<ShareCreationFlowStep, FlowError>;
async fn handle_check_manage_permission(
&self,
instance: &mut ShareCreationFlowInstance,
) -> Result<Option<ShareCreationFlowStep>, FlowError>;
async fn evaluate_check_share_type(
&self,
instance: &ShareCreationFlowInstance,
) -> Result<ShareCreationFlowStep, FlowError>;
async fn handle_validate_public_share(
&self,
instance: &mut ShareCreationFlowInstance,
) -> Result<Option<ShareCreationFlowStep>, FlowError>;
async fn handle_validate_user_share(
&self,
instance: &mut ShareCreationFlowInstance,
) -> Result<Option<ShareCreationFlowStep>, FlowError>;
async fn handle_validate_email_share(
&self,
instance: &mut ShareCreationFlowInstance,
) -> Result<Option<ShareCreationFlowStep>, FlowError>;
async fn handle_check_existing_share(
&self,
instance: &mut ShareCreationFlowInstance,
) -> Result<Option<ShareCreationFlowStep>, FlowError>;
async fn handle_update_existing_share(
&self,
instance: &mut ShareCreationFlowInstance,
) -> Result<Option<ShareCreationFlowStep>, FlowError>;
async fn handle_create_share_record(
&self,
instance: &mut ShareCreationFlowInstance,
) -> Result<Option<ShareCreationFlowStep>, FlowError>;
async fn handle_generate_share_url(
&self,
instance: &mut ShareCreationFlowInstance,
) -> Result<Option<ShareCreationFlowStep>, FlowError>;
async fn evaluate_determine_notification(
&self,
instance: &ShareCreationFlowInstance,
) -> Result<ShareCreationFlowStep, FlowError>;
async fn handle_notify_user(
&self,
instance: &mut ShareCreationFlowInstance,
) -> Result<Option<ShareCreationFlowStep>, FlowError>;
async fn handle_send_email_invitation(
&self,
instance: &mut ShareCreationFlowInstance,
) -> Result<Option<ShareCreationFlowStep>, FlowError>;
async fn handle_log_share(
&self,
instance: &mut ShareCreationFlowInstance,
) -> Result<Option<ShareCreationFlowStep>, FlowError>;
async fn handle_complete(
&self,
instance: &mut ShareCreationFlowInstance,
) -> Result<(), FlowError>;
async fn handle_file_not_found(
&self,
instance: &mut ShareCreationFlowInstance,
) -> Result<(), FlowError>;
async fn handle_permission_denied(
&self,
instance: &mut ShareCreationFlowInstance,
) -> Result<Option<ShareCreationFlowStep>, FlowError>;
async fn handle_failed_permission(
&self,
instance: &mut ShareCreationFlowInstance,
) -> Result<(), FlowError>;
async fn handle_public_sharing_disabled(
&self,
instance: &mut ShareCreationFlowInstance,
) -> Result<(), FlowError>;
async fn handle_user_not_found(
&self,
instance: &mut ShareCreationFlowInstance,
) -> Result<(), FlowError>;
async fn handle_invalid_config(
&self,
instance: &mut ShareCreationFlowInstance,
) -> Result<(), FlowError>;
}
pub struct ShareCreationFlowExecutor<H: ShareCreationStepHandler> {
handler: Arc<H>,
}
impl<H: ShareCreationStepHandler> ShareCreationFlowExecutor<H> {
pub fn new(handler: Arc<H>) -> Self {
Self { handler }
}
pub async fn start(&self, instance_id: impl Into<String>) -> Result<ShareCreationFlowInstance, FlowError> {
let mut instance = ShareCreationFlowInstance::new(instance_id);
instance.status = ShareCreationFlowStatus::Running;
instance.current_step = Some(ShareCreationFlowStep::ValidateFile);
Ok(instance)
}
pub async fn execute_step(&self, instance: &mut ShareCreationFlowInstance) -> Result<(), FlowError> {
let current_step = match instance.current_step {
Some(step) => step,
None => return Err(FlowError::NoCurrentStep),
};
let next_step = match current_step {
ShareCreationFlowStep::ValidateFile => {
self.handler.handle_validate_file(instance).await?
}
ShareCreationFlowStep::CheckSharePermission => {
let next = self.handler.evaluate_check_share_permission(instance).await?;
Some(next)
}
ShareCreationFlowStep::CheckManagePermission => {
self.handler.handle_check_manage_permission(instance).await?
}
ShareCreationFlowStep::CheckShareType => {
let next = self.handler.evaluate_check_share_type(instance).await?;
Some(next)
}
ShareCreationFlowStep::ValidatePublicShare => {
self.handler.handle_validate_public_share(instance).await?
}
ShareCreationFlowStep::ValidateUserShare => {
self.handler.handle_validate_user_share(instance).await?
}
ShareCreationFlowStep::ValidateEmailShare => {
self.handler.handle_validate_email_share(instance).await?
}
ShareCreationFlowStep::CheckExistingShare => {
self.handler.handle_check_existing_share(instance).await?
}
ShareCreationFlowStep::UpdateExistingShare => {
self.handler.handle_update_existing_share(instance).await?
}
ShareCreationFlowStep::CreateShareRecord => {
self.handler.handle_create_share_record(instance).await?
}
ShareCreationFlowStep::GenerateShareUrl => {
self.handler.handle_generate_share_url(instance).await?
}
ShareCreationFlowStep::DetermineNotification => {
let next = self.handler.evaluate_determine_notification(instance).await?;
Some(next)
}
ShareCreationFlowStep::NotifyUser => {
self.handler.handle_notify_user(instance).await?
}
ShareCreationFlowStep::SendEmailInvitation => {
self.handler.handle_send_email_invitation(instance).await?
}
ShareCreationFlowStep::LogShare => {
self.handler.handle_log_share(instance).await?
}
ShareCreationFlowStep::Complete => {
self.handler.handle_complete(instance).await?;
None }
ShareCreationFlowStep::FileNotFound => {
self.handler.handle_file_not_found(instance).await?;
None }
ShareCreationFlowStep::PermissionDenied => {
self.handler.handle_permission_denied(instance).await?
}
ShareCreationFlowStep::FailedPermission => {
self.handler.handle_failed_permission(instance).await?;
None }
ShareCreationFlowStep::PublicSharingDisabled => {
self.handler.handle_public_sharing_disabled(instance).await?;
None }
ShareCreationFlowStep::UserNotFound => {
self.handler.handle_user_not_found(instance).await?;
None }
ShareCreationFlowStep::InvalidConfig => {
self.handler.handle_invalid_config(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 = ShareCreationFlowStatus::Completed;
}
}
instance.updated_at = chrono::Utc::now();
Ok(())
}
pub async fn run(&self, instance: &mut ShareCreationFlowInstance) -> Result<(), FlowError> {
while !instance.is_complete() && instance.status != ShareCreationFlowStatus::Waiting {
self.execute_step(instance).await?;
}
Ok(())
}
pub fn cancel(&self, instance: &mut ShareCreationFlowInstance) {
instance.status = ShareCreationFlowStatus::Cancelled;
instance.updated_at = chrono::Utc::now();
}
pub fn fail(&self, instance: &mut ShareCreationFlowInstance, error: impl Into<String>) {
instance.status = ShareCreationFlowStatus::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 = ShareCreationFlowInstance::new("test-1");
assert_eq!(instance.id, "test-1");
assert_eq!(instance.status, ShareCreationFlowStatus::Pending);
assert!(!instance.is_complete());
}
#[test]
fn test_flow_context() {
let mut instance = ShareCreationFlowInstance::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 = ShareCreationFlowInstance::new("test-3");
assert!(!instance.is_running());
instance.status = ShareCreationFlowStatus::Running;
assert!(instance.is_running());
assert!(!instance.is_complete());
instance.status = ShareCreationFlowStatus::Completed;
assert!(instance.is_complete());
assert!(!instance.is_running());
}
}