use std::sync::Arc;
#[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, serde::Serialize, serde::Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum FlowStatus {
Pending,
Running,
Waiting,
Completed,
Failed,
Cancelled,
Compensating,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(bound(
serialize = "S: serde::Serialize",
deserialize = "S: serde::de::DeserializeOwned",
))]
pub struct FlowInstance<S>
where
S: std::fmt::Debug + Clone + serde::Serialize + serde::de::DeserializeOwned,
{
pub id: String,
pub status: FlowStatus,
pub current_step: Option<S>,
pub context: serde_json::Value,
pub completed_steps: Vec<S>,
pub error: Option<String>,
pub created_at: chrono::DateTime<chrono::Utc>,
pub updated_at: chrono::DateTime<chrono::Utc>,
}
impl<S> FlowInstance<S>
where
S: std::fmt::Debug + Clone + serde::Serialize + serde::de::DeserializeOwned,
{
pub fn new(id: impl Into<String>) -> Self {
let now = chrono::Utc::now();
Self {
id: id.into(),
status: FlowStatus::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,
FlowStatus::Completed | FlowStatus::Failed | FlowStatus::Cancelled
)
}
pub fn is_running(&self) -> bool {
matches!(self.status, FlowStatus::Running | FlowStatus::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
}
}
}
impl<S> WorkflowContext<FlowInstance<S>> for FlowInstance<S>
where
S: std::fmt::Debug + Clone + serde::Serialize + serde::de::DeserializeOwned + Send + Sync,
{
fn entity(&self) -> &FlowInstance<S> { 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) }
}
pub struct FlowExecutor<H> {
pub handler: Arc<H>,
}
impl<H: Send + Sync + 'static> FlowExecutor<H> {
pub fn new(handler: Arc<H>) -> Self {
Self { handler }
}
pub fn start<S>(&self, id: impl Into<String>, first_step: S) -> FlowInstance<S>
where
S: std::fmt::Debug + Clone + serde::Serialize + serde::de::DeserializeOwned,
{
let mut instance = FlowInstance::new(id);
instance.status = FlowStatus::Running;
instance.current_step = Some(first_step);
instance
}
pub fn cancel<S>(&self, instance: &mut FlowInstance<S>)
where
S: std::fmt::Debug + Clone + serde::Serialize + serde::de::DeserializeOwned,
{
instance.status = FlowStatus::Cancelled;
instance.updated_at = chrono::Utc::now();
}
pub fn fail<S>(&self, instance: &mut FlowInstance<S>, error: impl Into<String>)
where
S: std::fmt::Debug + Clone + serde::Serialize + serde::de::DeserializeOwned,
{
instance.status = FlowStatus::Failed;
instance.error = Some(error.into());
instance.updated_at = chrono::Utc::now();
}
}
pub trait WorkflowContext<E>: Send + Sync {
fn entity(&self) -> &E;
fn set_var(&mut self, key: &str, value: serde_json::Value);
fn get_var(&self, key: &str) -> Option<&serde_json::Value>;
}
pub trait WorkflowStep: Send + Sync {
fn name(&self) -> &'static str;
}