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}