Skip to main content

backbone_core/
flow.rs

1//! Workflow step and context traits for saga-pattern workflows.
2//!
3//! Phase 0 generic base for the `flow.rs` generator (Category C).
4//!
5//! Every generated `{Name}Workflow` used to re-define the same
6//! structural boilerplate. This module provides the single generic version so
7//! that generated files collapse to type aliases:
8//!
9//! ```rust,ignore
10//! use backbone_core::flow::{
11//!     FlowError, FlowStatus, FlowInstance, FlowExecutor,
12//!     WorkflowStep, WorkflowContext,
13//! };
14//! use crate::application::workflows::order_processing_workflow::OrderProcessingFlowStep;
15//!
16//! pub type OrderProcessingFlowError    = FlowError;
17//! pub type OrderProcessingFlowStatus   = FlowStatus;
18//! pub type OrderProcessingFlowInstance = FlowInstance<OrderProcessingFlowStep>;
19//!
20//! pub struct OrderProcessingFlowExecutor<H>(FlowExecutor<H>);
21//! // impl execute_step() — the only unique part — stays in the generated file
22//! ```
23
24use std::sync::Arc;
25
26// ─── FlowError ───────────────────────────────────────────────────────────────
27
28/// Common workflow execution error — identical for every workflow.
29///
30/// Generated files use a type alias: `pub type OrderProcessingFlowError = FlowError;`
31#[derive(Debug, Clone, thiserror::Error)]
32pub enum FlowError {
33    #[error("No current step to execute")]
34    NoCurrentStep,
35
36    #[error("Step execution failed: {0}")]
37    StepFailed(String),
38
39    #[error("Condition evaluation failed: {0}")]
40    ConditionFailed(String),
41
42    #[error("Compensation failed: {0}")]
43    CompensationFailed(String),
44
45    #[error("Flow timed out")]
46    Timeout,
47
48    #[error("Flow cancelled")]
49    Cancelled,
50
51    #[error("Invalid state transition: {from} -> {to}")]
52    InvalidTransition { from: String, to: String },
53}
54
55// ─── FlowStatus ──────────────────────────────────────────────────────────────
56
57/// Workflow execution status — identical for every workflow.
58///
59/// Generated files use a type alias: `pub type OrderProcessingFlowStatus = FlowStatus;`
60#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
61#[serde(rename_all = "snake_case")]
62pub enum FlowStatus {
63    /// Flow is pending execution
64    Pending,
65    /// Flow is currently running
66    Running,
67    /// Flow is waiting for an event or condition
68    Waiting,
69    /// Flow completed successfully
70    Completed,
71    /// Flow failed
72    Failed,
73    /// Flow was cancelled
74    Cancelled,
75    /// Flow is compensating (rolling back)
76    Compensating,
77}
78
79// ─── FlowInstance ─────────────────────────────────────────────────────────────
80
81/// Generic workflow instance parametrised over the step enum `S`.
82///
83/// Generated files use a type alias:
84/// `pub type OrderProcessingFlowInstance = FlowInstance<OrderProcessingFlowStep>;`
85#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
86#[serde(bound(
87    serialize = "S: serde::Serialize",
88    deserialize = "S: serde::de::DeserializeOwned",
89))]
90pub struct FlowInstance<S>
91where
92    S: std::fmt::Debug + Clone + serde::Serialize + serde::de::DeserializeOwned,
93{
94    pub id: String,
95    pub status: FlowStatus,
96    pub current_step: Option<S>,
97    pub context: serde_json::Value,
98    pub completed_steps: Vec<S>,
99    pub error: Option<String>,
100    pub created_at: chrono::DateTime<chrono::Utc>,
101    pub updated_at: chrono::DateTime<chrono::Utc>,
102}
103
104impl<S> FlowInstance<S>
105where
106    S: std::fmt::Debug + Clone + serde::Serialize + serde::de::DeserializeOwned,
107{
108    pub fn new(id: impl Into<String>) -> Self {
109        let now = chrono::Utc::now();
110        Self {
111            id: id.into(),
112            status: FlowStatus::Pending,
113            current_step: None,
114            context: serde_json::json!({}),
115            completed_steps: Vec::new(),
116            error: None,
117            created_at: now,
118            updated_at: now,
119        }
120    }
121
122    pub fn is_complete(&self) -> bool {
123        matches!(
124            self.status,
125            FlowStatus::Completed | FlowStatus::Failed | FlowStatus::Cancelled
126        )
127    }
128
129    pub fn is_running(&self) -> bool {
130        matches!(self.status, FlowStatus::Running | FlowStatus::Waiting)
131    }
132
133    pub fn set_context(&mut self, key: &str, value: serde_json::Value) {
134        if let serde_json::Value::Object(ref mut map) = self.context {
135            map.insert(key.to_string(), value);
136        }
137        self.updated_at = chrono::Utc::now();
138    }
139
140    pub fn get_context(&self, key: &str) -> Option<&serde_json::Value> {
141        if let serde_json::Value::Object(ref map) = self.context {
142            map.get(key)
143        } else {
144            None
145        }
146    }
147}
148
149// ─── WorkflowContext impl on FlowInstance ────────────────────────────────────
150
151impl<S> WorkflowContext<FlowInstance<S>> for FlowInstance<S>
152where
153    S: std::fmt::Debug + Clone + serde::Serialize + serde::de::DeserializeOwned + Send + Sync,
154{
155    fn entity(&self) -> &FlowInstance<S> { self }
156    fn set_var(&mut self, key: &str, value: serde_json::Value) { self.set_context(key, value); }
157    fn get_var(&self, key: &str) -> Option<&serde_json::Value> { self.get_context(key) }
158}
159
160// ─── FlowExecutor ────────────────────────────────────────────────────────────
161
162/// Generic executor shell — holds the handler, provides cancel/fail helpers.
163///
164/// Generated files wrap this with a newtype and add the unique `execute_step()`:
165///
166/// ```rust,ignore
167/// pub struct OrderProcessingFlowExecutor<H: OrderProcessingStepHandler>(FlowExecutor<H>);
168///
169/// impl<H: OrderProcessingStepHandler> OrderProcessingFlowExecutor<H> {
170///     pub fn new(handler: Arc<H>) -> Self { Self(FlowExecutor::new(handler)) }
171///     pub fn handler(&self) -> &Arc<H> { &self.0.handler }
172///
173///     pub async fn execute_step(&self, instance: &mut OrderProcessingFlowInstance)
174///         -> Result<(), FlowError> { /* unique dispatch */ }
175///
176///     pub async fn start(&self, id: impl Into<String>)
177///         -> Result<OrderProcessingFlowInstance, FlowError> { self.0.start(id) }
178///
179///     pub async fn run(&self, instance: &mut OrderProcessingFlowInstance)
180///         -> Result<(), FlowError> {
181///         while !instance.is_complete() && instance.status != FlowStatus::Waiting {
182///             self.execute_step(instance).await?;
183///         }
184///         Ok(())
185///     }
186///
187///     pub fn cancel(&self, instance: &mut OrderProcessingFlowInstance) {
188///         self.0.cancel(instance);
189///     }
190///     pub fn fail(&self, instance: &mut OrderProcessingFlowInstance, error: impl Into<String>) {
191///         self.0.fail(instance, error);
192///     }
193/// }
194/// ```
195pub struct FlowExecutor<H> {
196    pub handler: Arc<H>,
197}
198
199impl<H: Send + Sync + 'static> FlowExecutor<H> {
200    pub fn new(handler: Arc<H>) -> Self {
201        Self { handler }
202    }
203
204    /// Start a new flow instance at a given step.
205    pub fn start<S>(&self, id: impl Into<String>, first_step: S) -> FlowInstance<S>
206    where
207        S: std::fmt::Debug + Clone + serde::Serialize + serde::de::DeserializeOwned,
208    {
209        let mut instance = FlowInstance::new(id);
210        instance.status = FlowStatus::Running;
211        instance.current_step = Some(first_step);
212        instance
213    }
214
215    /// Cancel the flow instance.
216    pub fn cancel<S>(&self, instance: &mut FlowInstance<S>)
217    where
218        S: std::fmt::Debug + Clone + serde::Serialize + serde::de::DeserializeOwned,
219    {
220        instance.status = FlowStatus::Cancelled;
221        instance.updated_at = chrono::Utc::now();
222    }
223
224    /// Mark the flow instance as failed with a message.
225    pub fn fail<S>(&self, instance: &mut FlowInstance<S>, error: impl Into<String>)
226    where
227        S: std::fmt::Debug + Clone + serde::Serialize + serde::de::DeserializeOwned,
228    {
229        instance.status = FlowStatus::Failed;
230        instance.error = Some(error.into());
231        instance.updated_at = chrono::Utc::now();
232    }
233}
234
235// ─── Traits ──────────────────────────────────────────────────────────────────
236
237/// Workflow execution context for an entity `E`.
238pub trait WorkflowContext<E>: Send + Sync {
239    fn entity(&self) -> &E;
240    fn set_var(&mut self, key: &str, value: serde_json::Value);
241    fn get_var(&self, key: &str) -> Option<&serde_json::Value>;
242}
243
244/// A named step within a workflow.
245pub trait WorkflowStep: Send + Sync {
246    fn name(&self) -> &'static str;
247}