use crate::errors::CoreError;
use serde::{Deserialize, Serialize};
use serde_json::Value;
use std::time::{SystemTime, UNIX_EPOCH};
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum JobStatus {
Pending,
Queued,
Running,
Paused,
Succeeded,
Failed,
Cancelled,
}
impl JobStatus {
pub fn is_terminal(&self) -> bool {
matches!(
self,
JobStatus::Succeeded | JobStatus::Failed | JobStatus::Cancelled
)
}
pub fn is_success(&self) -> bool {
*self == JobStatus::Succeeded
}
pub fn from_str(s: &str) -> Option<Self> {
let normalized = s.trim().to_lowercase().replace(' ', "_");
match normalized.as_str() {
"pending" => Some(JobStatus::Pending),
"queued" => Some(JobStatus::Queued),
"running" | "in_progress" => Some(JobStatus::Running),
"paused" => Some(JobStatus::Paused),
"succeeded" | "success" | "completed" | "complete" => Some(JobStatus::Succeeded),
"failed" | "failure" | "error" => Some(JobStatus::Failed),
"cancelled" | "canceled" | "cancel" => Some(JobStatus::Cancelled),
_ => None,
}
}
pub fn as_str(&self) -> &'static str {
match self {
JobStatus::Pending => "pending",
JobStatus::Queued => "queued",
JobStatus::Running => "running",
JobStatus::Paused => "paused",
JobStatus::Succeeded => "succeeded",
JobStatus::Failed => "failed",
JobStatus::Cancelled => "cancelled",
}
}
}
impl std::fmt::Display for JobStatus {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{}", self.as_str())
}
}
impl Default for JobStatus {
fn default() -> Self {
JobStatus::Pending
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum CandidateStatus {
Evaluating,
Completed,
Accepted,
Rejected,
Failed,
}
impl CandidateStatus {
pub fn is_terminal(&self) -> bool {
matches!(
self,
CandidateStatus::Completed
| CandidateStatus::Accepted
| CandidateStatus::Rejected
| CandidateStatus::Failed
)
}
pub fn as_str(&self) -> &'static str {
match self {
CandidateStatus::Evaluating => "evaluating",
CandidateStatus::Completed => "completed",
CandidateStatus::Accepted => "accepted",
CandidateStatus::Rejected => "rejected",
CandidateStatus::Failed => "failed",
}
}
}
impl std::fmt::Display for CandidateStatus {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{}", self.as_str())
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
pub enum JobEventType {
#[serde(rename = "job.created")]
JobCreated,
#[serde(rename = "job.queued")]
JobQueued,
#[serde(rename = "job.in_progress")]
JobInProgress,
#[serde(rename = "job.completed")]
JobCompleted,
#[serde(rename = "job.failed")]
JobFailed,
#[serde(rename = "job.cancelled")]
JobCancelled,
}
impl JobEventType {
pub fn as_str(&self) -> &'static str {
match self {
JobEventType::JobCreated => "job.created",
JobEventType::JobQueued => "job.queued",
JobEventType::JobInProgress => "job.in_progress",
JobEventType::JobCompleted => "job.completed",
JobEventType::JobFailed => "job.failed",
JobEventType::JobCancelled => "job.cancelled",
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct JobEvent {
#[serde(rename = "type")]
pub event_type: String,
pub job_id: String,
pub seq: i64,
pub timestamp: f64,
#[serde(skip_serializing_if = "Option::is_none")]
pub data: Option<Value>,
#[serde(skip_serializing_if = "Option::is_none")]
pub message: Option<String>,
}
#[derive(Debug, Clone)]
pub struct JobLifecycle {
job_id: String,
status: JobStatus,
events: Vec<JobEvent>,
started_at: Option<f64>,
ended_at: Option<f64>,
}
impl JobLifecycle {
pub fn new(job_id: &str) -> Self {
Self {
job_id: job_id.to_string(),
status: JobStatus::Pending,
events: Vec::new(),
started_at: None,
ended_at: None,
}
}
pub fn status(&self) -> JobStatus {
self.status
}
pub fn job_id(&self) -> &str {
&self.job_id
}
pub fn elapsed_seconds(&self) -> Option<f64> {
let start = self.started_at?;
let end = self.ended_at.unwrap_or_else(now_timestamp);
Some(end - start)
}
pub fn events(&self) -> &[JobEvent] {
&self.events
}
fn emit(
&mut self,
event_type: JobEventType,
data: Option<Value>,
message: Option<&str>,
) -> JobEvent {
let event = JobEvent {
event_type: event_type.as_str().to_string(),
job_id: self.job_id.clone(),
seq: (self.events.len() + 1) as i64,
timestamp: now_timestamp(),
data,
message: message.map(String::from),
};
self.events.push(event.clone());
event
}
pub fn start(&mut self) -> Result<JobEvent, CoreError> {
self.start_with_data(None, None)
}
pub fn start_with_data(
&mut self,
data: Option<Value>,
message: Option<&str>,
) -> Result<JobEvent, CoreError> {
if self.status != JobStatus::Pending {
return Err(CoreError::Job(crate::errors::JobErrorInfo {
job_id: self.job_id.clone(),
message: format!("Cannot start job in {} status", self.status),
code: Some("INVALID_TRANSITION".to_string()),
}));
}
self.status = JobStatus::Running;
self.started_at = Some(now_timestamp());
Ok(self.emit(
JobEventType::JobInProgress,
data,
Some(message.unwrap_or("Job started")),
))
}
pub fn complete(&mut self, data: Option<Value>) -> Result<JobEvent, CoreError> {
self.complete_with_message(data, None)
}
pub fn complete_with_message(
&mut self,
data: Option<Value>,
message: Option<&str>,
) -> Result<JobEvent, CoreError> {
if self.status != JobStatus::Running {
return Err(CoreError::Job(crate::errors::JobErrorInfo {
job_id: self.job_id.clone(),
message: format!("Cannot complete job in {} status", self.status),
code: Some("INVALID_TRANSITION".to_string()),
}));
}
self.status = JobStatus::Succeeded;
self.ended_at = Some(now_timestamp());
let mut event_data = data.unwrap_or_else(|| Value::Object(Default::default()));
if let Value::Object(ref mut map) = event_data {
if let Some(elapsed) = self.elapsed_seconds() {
map.insert("elapsed_seconds".to_string(), Value::from(elapsed));
}
}
Ok(self.emit(
JobEventType::JobCompleted,
Some(event_data),
Some(message.unwrap_or("Job completed successfully")),
))
}
pub fn fail(&mut self, error: Option<&str>) -> Result<JobEvent, CoreError> {
self.fail_with_data(error, None)
}
pub fn fail_with_data(
&mut self,
error: Option<&str>,
data: Option<Value>,
) -> Result<JobEvent, CoreError> {
if self.status != JobStatus::Running && self.status != JobStatus::Pending {
return Err(CoreError::Job(crate::errors::JobErrorInfo {
job_id: self.job_id.clone(),
message: format!("Cannot fail job in {} status", self.status),
code: Some("INVALID_TRANSITION".to_string()),
}));
}
self.status = JobStatus::Failed;
self.ended_at = Some(now_timestamp());
let mut event_data = data.unwrap_or_else(|| Value::Object(Default::default()));
if let Value::Object(ref mut map) = event_data {
if let Some(err) = error {
map.insert("error".to_string(), Value::String(err.to_string()));
}
if let Some(elapsed) = self.elapsed_seconds() {
map.insert("elapsed_seconds".to_string(), Value::from(elapsed));
}
}
Ok(self.emit(
JobEventType::JobFailed,
Some(event_data),
Some(error.unwrap_or("Job failed")),
))
}
pub fn cancel(&mut self) -> Result<JobEvent, CoreError> {
self.cancel_with_message(None)
}
pub fn cancel_with_message(&mut self, message: Option<&str>) -> Result<JobEvent, CoreError> {
if self.status.is_terminal() {
return Err(CoreError::Job(crate::errors::JobErrorInfo {
job_id: self.job_id.clone(),
message: format!("Cannot cancel job in {} status", self.status),
code: Some("INVALID_TRANSITION".to_string()),
}));
}
self.status = JobStatus::Cancelled;
self.ended_at = Some(now_timestamp());
let mut event_data = Value::Object(Default::default());
if let Value::Object(ref mut map) = event_data {
if let Some(elapsed) = self.elapsed_seconds() {
map.insert("elapsed_seconds".to_string(), Value::from(elapsed));
}
}
Ok(self.emit(
JobEventType::JobCancelled,
Some(event_data),
Some(message.unwrap_or("Job cancelled")),
))
}
}
fn now_timestamp() -> f64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|d| d.as_secs_f64())
.unwrap_or(0.0)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_job_status_from_str() {
assert_eq!(JobStatus::from_str("pending"), Some(JobStatus::Pending));
assert_eq!(JobStatus::from_str("RUNNING"), Some(JobStatus::Running));
assert_eq!(JobStatus::from_str("in_progress"), Some(JobStatus::Running));
assert_eq!(JobStatus::from_str("paused"), Some(JobStatus::Paused));
assert_eq!(JobStatus::from_str("success"), Some(JobStatus::Succeeded));
assert_eq!(JobStatus::from_str("completed"), Some(JobStatus::Succeeded));
assert_eq!(JobStatus::from_str("failed"), Some(JobStatus::Failed));
assert_eq!(JobStatus::from_str("cancelled"), Some(JobStatus::Cancelled));
assert_eq!(JobStatus::from_str("canceled"), Some(JobStatus::Cancelled));
assert_eq!(JobStatus::from_str("unknown"), None);
}
#[test]
fn test_job_status_is_terminal() {
assert!(!JobStatus::Pending.is_terminal());
assert!(!JobStatus::Queued.is_terminal());
assert!(!JobStatus::Running.is_terminal());
assert!(!JobStatus::Paused.is_terminal());
assert!(JobStatus::Succeeded.is_terminal());
assert!(JobStatus::Failed.is_terminal());
assert!(JobStatus::Cancelled.is_terminal());
}
#[test]
fn test_job_lifecycle_happy_path() {
let mut lifecycle = JobLifecycle::new("test-job-123");
assert_eq!(lifecycle.status(), JobStatus::Pending);
let start_event = lifecycle.start().unwrap();
assert_eq!(lifecycle.status(), JobStatus::Running);
assert_eq!(start_event.event_type, "job.in_progress");
assert_eq!(start_event.seq, 1);
let complete_event = lifecycle.complete(None).unwrap();
assert_eq!(lifecycle.status(), JobStatus::Succeeded);
assert_eq!(complete_event.event_type, "job.completed");
assert_eq!(complete_event.seq, 2);
assert!(lifecycle.elapsed_seconds().is_some());
}
#[test]
fn test_job_lifecycle_fail() {
let mut lifecycle = JobLifecycle::new("test-job-456");
lifecycle.start().unwrap();
let fail_event = lifecycle.fail(Some("Something went wrong")).unwrap();
assert_eq!(lifecycle.status(), JobStatus::Failed);
assert_eq!(fail_event.event_type, "job.failed");
if let Some(data) = &fail_event.data {
assert!(data.get("error").is_some());
}
}
#[test]
fn test_job_lifecycle_cancel() {
let mut lifecycle = JobLifecycle::new("test-job-789");
lifecycle.start().unwrap();
let cancel_event = lifecycle.cancel().unwrap();
assert_eq!(lifecycle.status(), JobStatus::Cancelled);
assert_eq!(cancel_event.event_type, "job.cancelled");
}
#[test]
fn test_invalid_transitions() {
let mut lifecycle = JobLifecycle::new("test-job-invalid");
assert!(lifecycle.complete(None).is_err());
lifecycle.start().unwrap();
assert!(lifecycle.start().is_err());
lifecycle.complete(None).unwrap();
assert!(lifecycle.fail(None).is_err());
assert!(lifecycle.cancel().is_err());
}
}