use serde::{Deserialize, Serialize};
use std::str::FromStr;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum ProcessingJobState {
Pending,
Running,
Completed,
Failed,
Cancelled,
}
impl Default for ProcessingJobState {
fn default() -> Self {
Self::Pending
}
}
impl std::fmt::Display for ProcessingJobState {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::Pending => write!(f, "pending"),
Self::Running => write!(f, "running"),
Self::Completed => write!(f, "completed"),
Self::Failed => write!(f, "failed"),
Self::Cancelled => write!(f, "cancelled"),
}
}
}
impl FromStr for ProcessingJobState {
type Err = StateMachineError;
fn from_str(s: &str) -> Result<Self, Self::Err> {
match s.to_lowercase().as_str() {
"pending" => Ok(Self::Pending),
"running" => Ok(Self::Running),
"completed" => Ok(Self::Completed),
"failed" => Ok(Self::Failed),
"cancelled" => Ok(Self::Cancelled),
_ => Err(StateMachineError::InvalidState(s.to_string())),
}
}
}
impl ProcessingJobState {
pub fn is_initial(&self) -> bool {
matches!(self, Self::Pending)
}
pub fn is_final(&self) -> bool {
matches!(self, Self::Completed | Self::Failed | Self::Cancelled)
}
pub fn all() -> Vec<Self> {
vec![
Self::Pending,
Self::Running,
Self::Completed,
Self::Failed,
Self::Cancelled,
]
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum ProcessingJobTransition {
Start,
Complete,
Fail,
CancelPending,
CancelRunning,
Retry,
Reprocess,
}
impl std::fmt::Display for ProcessingJobTransition {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::Start => write!(f, "start"),
Self::Complete => write!(f, "complete"),
Self::Fail => write!(f, "fail"),
Self::CancelPending => write!(f, "cancel_pending"),
Self::CancelRunning => write!(f, "cancel_running"),
Self::Retry => write!(f, "retry"),
Self::Reprocess => write!(f, "reprocess"),
}
}
}
impl FromStr for ProcessingJobTransition {
type Err = StateMachineError;
fn from_str(s: &str) -> Result<Self, Self::Err> {
match s.to_lowercase().as_str() {
"start" => Ok(Self::Start),
"complete" => Ok(Self::Complete),
"fail" => Ok(Self::Fail),
"cancel_pending" => Ok(Self::CancelPending),
"cancel_running" => Ok(Self::CancelRunning),
"retry" => Ok(Self::Retry),
"reprocess" => Ok(Self::Reprocess),
_ => Err(StateMachineError::InvalidTransition(s.to_string())),
}
}
}
impl ProcessingJobTransition {
pub fn target_state(&self) -> ProcessingJobState {
match self {
Self::Start => ProcessingJobState::Running,
Self::Complete => ProcessingJobState::Completed,
Self::Fail => ProcessingJobState::Failed,
Self::CancelPending => ProcessingJobState::Cancelled,
Self::CancelRunning => ProcessingJobState::Cancelled,
Self::Retry => ProcessingJobState::Pending,
Self::Reprocess => ProcessingJobState::Pending,
}
}
pub fn all() -> Vec<Self> {
vec![
Self::Start,
Self::Complete,
Self::Fail,
Self::CancelPending,
Self::CancelRunning,
Self::Retry,
Self::Reprocess,
]
}
pub fn allowed_roles(&self) -> &'static [&'static str] {
match self {
Self::Start => &["system"],
Self::Complete => &["system"],
Self::Fail => &["system"],
Self::CancelPending => &["owner", "admin"],
Self::CancelRunning => &["owner", "admin"],
Self::Retry => &["system", "owner"],
Self::Reprocess => &["owner", "admin"],
}
}
}
use super::StateMachineError;
#[derive(Debug, Clone)]
pub struct ProcessingJobStateMachine {
current_state: ProcessingJobState,
}
impl ProcessingJobStateMachine {
pub fn new() -> Self {
Self {
current_state: ProcessingJobState::default(),
}
}
pub fn from_state(state: ProcessingJobState) -> Self {
Self { current_state: state }
}
pub fn current_state(&self) -> ProcessingJobState {
self.current_state
}
pub fn can_transition(&self, transition: ProcessingJobTransition) -> bool {
if matches!(self.current_state, ProcessingJobState::Cancelled) {
return false;
}
match (self.current_state, transition) {
(ProcessingJobState::Pending, ProcessingJobTransition::Start) => true,
(ProcessingJobState::Running, ProcessingJobTransition::Complete) => true,
(ProcessingJobState::Running, ProcessingJobTransition::Fail) => true,
(ProcessingJobState::Pending, ProcessingJobTransition::CancelPending) => true,
(ProcessingJobState::Running, ProcessingJobTransition::CancelRunning) => true,
(ProcessingJobState::Failed, ProcessingJobTransition::Retry) => true,
(ProcessingJobState::Completed, ProcessingJobTransition::Reprocess) => true,
_ => false,
}
}
pub fn can_transition_with_role(&self, transition: ProcessingJobTransition, role: &str) -> bool {
if !self.can_transition(transition) {
return false;
}
let allowed_roles = transition.allowed_roles();
if allowed_roles.is_empty() {
return true; }
allowed_roles.iter().any(|r| *r == role || *r == "*")
}
pub fn transition(&mut self, transition: ProcessingJobTransition) -> Result<ProcessingJobState, StateMachineError> {
if !self.can_transition(transition) {
return Err(StateMachineError::TransitionNotAllowed {
transition: transition.to_string(),
from: self.current_state.to_string(),
});
}
self.current_state = transition.target_state();
Ok(self.current_state)
}
pub fn transition_with_role(&mut self, transition: ProcessingJobTransition, role: &str) -> Result<ProcessingJobState, StateMachineError> {
if !self.can_transition(transition) {
return Err(StateMachineError::TransitionNotAllowed {
transition: transition.to_string(),
from: self.current_state.to_string(),
});
}
if !self.can_transition_with_role(transition, role) {
return Err(StateMachineError::RoleNotAuthorized {
role: role.to_string(),
transition: transition.to_string(),
});
}
self.current_state = transition.target_state();
Ok(self.current_state)
}
pub fn available_transitions(&self) -> Vec<ProcessingJobTransition> {
ProcessingJobTransition::all()
.into_iter()
.filter(|t| self.can_transition(*t))
.collect()
}
pub fn available_transitions_for_role(&self, role: &str) -> Vec<ProcessingJobTransition> {
ProcessingJobTransition::all()
.into_iter()
.filter(|t| self.can_transition_with_role(*t, role))
.collect()
}
pub fn transition_to_state(&mut self, target: ProcessingJobState) -> Result<ProcessingJobState, StateMachineError> {
let valid = ProcessingJobTransition::all().into_iter()
.filter(|t| self.can_transition(*t))
.find(|t| t.target_state() == target);
match valid {
Some(t) => self.transition(t),
None => Err(StateMachineError::TransitionNotAllowed {
transition: target.to_string(),
from: self.current_state.to_string(),
}),
}
}
}
impl Default for ProcessingJobStateMachine {
fn default() -> Self {
Self::new()
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_initial_state() {
let sm = ProcessingJobStateMachine::new();
assert_eq!(sm.current_state(), ProcessingJobState::Pending);
assert!(sm.current_state().is_initial());
}
#[test]
fn test_valid_transition() {
let mut sm = ProcessingJobStateMachine::from_state(ProcessingJobState::Pending);
assert!(sm.can_transition(ProcessingJobTransition::Start));
let result = sm.transition(ProcessingJobTransition::Start);
assert!(result.is_ok());
assert_eq!(sm.current_state(), ProcessingJobState::Running);
}
#[test]
fn test_invalid_transition() {
let mut sm = ProcessingJobStateMachine::from_state(ProcessingJobState::Pending);
let result = sm.transition(ProcessingJobTransition::Complete);
assert!(result.is_err());
}
#[test]
fn test_state_parsing() {
let state: ProcessingJobState = "pending".parse().unwrap();
assert_eq!(state, ProcessingJobState::Pending);
}
#[test]
fn test_available_transitions() {
let sm = ProcessingJobStateMachine::new();
let available = sm.available_transitions();
assert!(!available.is_empty() || sm.current_state().is_final());
}
}