backbone_bucket/domain/state_machine/
processing_job_state_machine.rs1use serde::{Deserialize, Serialize};
6use std::str::FromStr;
7
8#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
10#[serde(rename_all = "snake_case")]
11pub enum ProcessingJobState {
12 Pending,
14 Running,
15 Completed,
17 Failed,
19 Cancelled,
21}
22
23impl Default for ProcessingJobState {
24 fn default() -> Self {
25 Self::Pending
26 }
27}
28
29impl std::fmt::Display for ProcessingJobState {
30 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
31 match self {
32 Self::Pending => write!(f, "pending"),
33 Self::Running => write!(f, "running"),
34 Self::Completed => write!(f, "completed"),
35 Self::Failed => write!(f, "failed"),
36 Self::Cancelled => write!(f, "cancelled"),
37 }
38 }
39}
40
41impl FromStr for ProcessingJobState {
42 type Err = StateMachineError;
43
44 fn from_str(s: &str) -> Result<Self, Self::Err> {
45 match s.to_lowercase().as_str() {
46 "pending" => Ok(Self::Pending),
47 "running" => Ok(Self::Running),
48 "completed" => Ok(Self::Completed),
49 "failed" => Ok(Self::Failed),
50 "cancelled" => Ok(Self::Cancelled),
51 _ => Err(StateMachineError::InvalidState(s.to_string())),
52 }
53 }
54}
55
56impl ProcessingJobState {
57 pub fn is_initial(&self) -> bool {
59 matches!(self, Self::Pending)
60 }
61
62 pub fn is_final(&self) -> bool {
64 matches!(self, Self::Completed | Self::Failed | Self::Cancelled)
65 }
66
67 pub fn all() -> Vec<Self> {
69 vec![
70 Self::Pending,
71 Self::Running,
72 Self::Completed,
73 Self::Failed,
74 Self::Cancelled,
75 ]
76 }
77}
78
79#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
81#[serde(rename_all = "snake_case")]
82pub enum ProcessingJobTransition {
83 Start,
85 Complete,
87 Fail,
89 CancelPending,
91 CancelRunning,
93 Retry,
95 Reprocess,
97}
98
99impl std::fmt::Display for ProcessingJobTransition {
100 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
101 match self {
102 Self::Start => write!(f, "start"),
103 Self::Complete => write!(f, "complete"),
104 Self::Fail => write!(f, "fail"),
105 Self::CancelPending => write!(f, "cancel_pending"),
106 Self::CancelRunning => write!(f, "cancel_running"),
107 Self::Retry => write!(f, "retry"),
108 Self::Reprocess => write!(f, "reprocess"),
109 }
110 }
111}
112
113impl FromStr for ProcessingJobTransition {
114 type Err = StateMachineError;
115
116 fn from_str(s: &str) -> Result<Self, Self::Err> {
117 match s.to_lowercase().as_str() {
118 "start" => Ok(Self::Start),
119 "complete" => Ok(Self::Complete),
120 "fail" => Ok(Self::Fail),
121 "cancel_pending" => Ok(Self::CancelPending),
122 "cancel_running" => Ok(Self::CancelRunning),
123 "retry" => Ok(Self::Retry),
124 "reprocess" => Ok(Self::Reprocess),
125 _ => Err(StateMachineError::InvalidTransition(s.to_string())),
126 }
127 }
128}
129
130impl ProcessingJobTransition {
131 pub fn target_state(&self) -> ProcessingJobState {
133 match self {
134 Self::Start => ProcessingJobState::Running,
135 Self::Complete => ProcessingJobState::Completed,
136 Self::Fail => ProcessingJobState::Failed,
137 Self::CancelPending => ProcessingJobState::Cancelled,
138 Self::CancelRunning => ProcessingJobState::Cancelled,
139 Self::Retry => ProcessingJobState::Pending,
140 Self::Reprocess => ProcessingJobState::Pending,
141 }
142 }
143
144 pub fn all() -> Vec<Self> {
146 vec![
147 Self::Start,
148 Self::Complete,
149 Self::Fail,
150 Self::CancelPending,
151 Self::CancelRunning,
152 Self::Retry,
153 Self::Reprocess,
154 ]
155 }
156
157 pub fn allowed_roles(&self) -> &'static [&'static str] {
159 match self {
160 Self::Start => &["system"],
161 Self::Complete => &["system"],
162 Self::Fail => &["system"],
163 Self::CancelPending => &["owner", "admin"],
164 Self::CancelRunning => &["owner", "admin"],
165 Self::Retry => &["system", "owner"],
166 Self::Reprocess => &["owner", "admin"],
167 }
168 }
169}
170
171use super::StateMachineError;
172
173#[derive(Debug, Clone)]
175pub struct ProcessingJobStateMachine {
176 current_state: ProcessingJobState,
177}
178
179impl ProcessingJobStateMachine {
180 pub fn new() -> Self {
182 Self {
183 current_state: ProcessingJobState::default(),
184 }
185 }
186
187 pub fn from_state(state: ProcessingJobState) -> Self {
189 Self { current_state: state }
190 }
191
192 pub fn current_state(&self) -> ProcessingJobState {
194 self.current_state
195 }
196
197 pub fn can_transition(&self, transition: ProcessingJobTransition) -> bool {
199 if matches!(self.current_state, ProcessingJobState::Cancelled) {
200 return false;
201 }
202
203 match (self.current_state, transition) {
204 (ProcessingJobState::Pending, ProcessingJobTransition::Start) => true,
205 (ProcessingJobState::Running, ProcessingJobTransition::Complete) => true,
206 (ProcessingJobState::Running, ProcessingJobTransition::Fail) => true,
207 (ProcessingJobState::Pending, ProcessingJobTransition::CancelPending) => true,
208 (ProcessingJobState::Running, ProcessingJobTransition::CancelRunning) => true,
209 (ProcessingJobState::Failed, ProcessingJobTransition::Retry) => true,
210 (ProcessingJobState::Completed, ProcessingJobTransition::Reprocess) => true,
211 _ => false,
212 }
213 }
214
215 pub fn can_transition_with_role(&self, transition: ProcessingJobTransition, role: &str) -> bool {
217 if !self.can_transition(transition) {
218 return false;
219 }
220
221 let allowed_roles = transition.allowed_roles();
222 if allowed_roles.is_empty() {
223 return true; }
225
226 allowed_roles.iter().any(|r| *r == role || *r == "*")
227 }
228
229 pub fn transition(&mut self, transition: ProcessingJobTransition) -> Result<ProcessingJobState, StateMachineError> {
231 if !self.can_transition(transition) {
232 return Err(StateMachineError::TransitionNotAllowed {
233 transition: transition.to_string(),
234 from: self.current_state.to_string(),
235 });
236 }
237
238 self.current_state = transition.target_state();
239 Ok(self.current_state)
240 }
241
242 pub fn transition_with_role(&mut self, transition: ProcessingJobTransition, role: &str) -> Result<ProcessingJobState, StateMachineError> {
244 if !self.can_transition(transition) {
246 return Err(StateMachineError::TransitionNotAllowed {
247 transition: transition.to_string(),
248 from: self.current_state.to_string(),
249 });
250 }
251
252 if !self.can_transition_with_role(transition, role) {
254 return Err(StateMachineError::RoleNotAuthorized {
255 role: role.to_string(),
256 transition: transition.to_string(),
257 });
258 }
259
260 self.current_state = transition.target_state();
261 Ok(self.current_state)
262 }
263
264 pub fn available_transitions(&self) -> Vec<ProcessingJobTransition> {
266 ProcessingJobTransition::all()
267 .into_iter()
268 .filter(|t| self.can_transition(*t))
269 .collect()
270 }
271
272 pub fn available_transitions_for_role(&self, role: &str) -> Vec<ProcessingJobTransition> {
274 ProcessingJobTransition::all()
275 .into_iter()
276 .filter(|t| self.can_transition_with_role(*t, role))
277 .collect()
278 }
279
280 pub fn transition_to_state(&mut self, target: ProcessingJobState) -> Result<ProcessingJobState, StateMachineError> {
285 let valid = ProcessingJobTransition::all().into_iter()
286 .filter(|t| self.can_transition(*t))
287 .find(|t| t.target_state() == target);
288 match valid {
289 Some(t) => self.transition(t),
290 None => Err(StateMachineError::TransitionNotAllowed {
291 transition: target.to_string(),
292 from: self.current_state.to_string(),
293 }),
294 }
295 }
296}
297
298impl Default for ProcessingJobStateMachine {
299 fn default() -> Self {
300 Self::new()
301 }
302}
303
304#[cfg(test)]
305mod tests {
306 use super::*;
307
308 #[test]
309 fn test_initial_state() {
310 let sm = ProcessingJobStateMachine::new();
311 assert_eq!(sm.current_state(), ProcessingJobState::Pending);
312 assert!(sm.current_state().is_initial());
313 }
314
315 #[test]
316 fn test_valid_transition() {
317 let mut sm = ProcessingJobStateMachine::from_state(ProcessingJobState::Pending);
318 assert!(sm.can_transition(ProcessingJobTransition::Start));
319 let result = sm.transition(ProcessingJobTransition::Start);
320 assert!(result.is_ok());
321 assert_eq!(sm.current_state(), ProcessingJobState::Running);
322 }
323
324 #[test]
325 fn test_invalid_transition() {
326 let mut sm = ProcessingJobStateMachine::from_state(ProcessingJobState::Pending);
327 let result = sm.transition(ProcessingJobTransition::Complete);
329 assert!(result.is_err());
330 }
331
332 #[test]
333 fn test_state_parsing() {
334 let state: ProcessingJobState = "pending".parse().unwrap();
335 assert_eq!(state, ProcessingJobState::Pending);
336 }
337
338 #[test]
339 fn test_available_transitions() {
340 let sm = ProcessingJobStateMachine::new();
341 let available = sm.available_transitions();
342 assert!(!available.is_empty() || sm.current_state().is_final());
344 }
345}