structure
This commit is contained in:
parent
9fbcee0bea
commit
ae7ec3ea8a
13 changed files with 991 additions and 19 deletions
230
src/task.rs
Normal file
230
src/task.rs
Normal file
|
|
@ -0,0 +1,230 @@
|
|||
use chrono::{DateTime, Utc};
|
||||
use serde::{Deserialize, Serialize};
|
||||
use uuid::Uuid;
|
||||
|
||||
use crate::enums::{TaskResultType, TaskStatus};
|
||||
use crate::todo::ToDo;
|
||||
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
pub struct Task {
|
||||
pub id: Uuid,
|
||||
pub todo_id: Uuid,
|
||||
pub description: String,
|
||||
pub status: TaskStatus,
|
||||
pub assigned_krill: Option<Uuid>,
|
||||
pub assigned_pod: Option<Uuid>,
|
||||
pub result: Option<TaskResult>,
|
||||
pub logs: Vec<String>,
|
||||
pub started_at: Option<DateTime<Utc>>,
|
||||
pub completed_at: Option<DateTime<Utc>>,
|
||||
}
|
||||
|
||||
impl Default for Task {
|
||||
fn default() -> Self {
|
||||
Self {
|
||||
id: Uuid::new_v4(),
|
||||
todo_id: Uuid::nil(),
|
||||
description: String::new(),
|
||||
status: TaskStatus::default(),
|
||||
assigned_krill: None,
|
||||
assigned_pod: None,
|
||||
result: None,
|
||||
logs: Vec::new(),
|
||||
started_at: None,
|
||||
completed_at: None,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl Task {
|
||||
pub fn new(todo_id: Uuid, description: String) -> Self {
|
||||
Self {
|
||||
id: Uuid::new_v4(),
|
||||
todo_id,
|
||||
description,
|
||||
status: TaskStatus::default(),
|
||||
assigned_krill: None,
|
||||
assigned_pod: None,
|
||||
result: None,
|
||||
logs: Vec::new(),
|
||||
started_at: None,
|
||||
completed_at: None,
|
||||
}
|
||||
}
|
||||
|
||||
pub fn start(&mut self) -> Result<(), String> {
|
||||
match self.status {
|
||||
TaskStatus::Pending => {
|
||||
self.status = TaskStatus::Running;
|
||||
self.started_at = Some(Utc::now());
|
||||
Ok(())
|
||||
}
|
||||
_ => Err(format!("cannot start task from status: {:?}", self.status)),
|
||||
}
|
||||
}
|
||||
|
||||
pub fn complete(&mut self, result: TaskResult) -> Result<(), String> {
|
||||
match self.status {
|
||||
TaskStatus::Running => {
|
||||
self.status = TaskStatus::Completed;
|
||||
self.completed_at = Some(Utc::now());
|
||||
self.result = Some(result);
|
||||
Ok(())
|
||||
}
|
||||
_ => Err(format!(
|
||||
"cannot complete task from status: {:?}",
|
||||
self.status
|
||||
)),
|
||||
}
|
||||
}
|
||||
|
||||
pub fn fail(&mut self, error: String) -> Result<(), String> {
|
||||
match self.status {
|
||||
TaskStatus::Running => {
|
||||
self.status = TaskStatus::Failed;
|
||||
self.completed_at = Some(Utc::now());
|
||||
self.result = Some(TaskResult::Error(error));
|
||||
Ok(())
|
||||
}
|
||||
_ => Err(format!("cannot fail task from status: {:?}", self.status)),
|
||||
}
|
||||
}
|
||||
|
||||
pub fn retry(&mut self) -> Result<(), String> {
|
||||
match self.status {
|
||||
TaskStatus::Failed | TaskStatus::Pending => {
|
||||
self.status = TaskStatus::Pending;
|
||||
self.started_at = None;
|
||||
self.completed_at = None;
|
||||
self.result = None;
|
||||
self.logs.clear();
|
||||
Ok(())
|
||||
}
|
||||
_ => Err(format!("cannot retry task from status: {:?}", self.status)),
|
||||
}
|
||||
}
|
||||
|
||||
pub fn duration_ms(&self) -> Option<i64> {
|
||||
match (self.started_at, self.completed_at) {
|
||||
(Some(start), Some(end)) => Some((end - start).num_milliseconds()),
|
||||
_ => None,
|
||||
}
|
||||
}
|
||||
|
||||
pub fn add_log(&mut self, log: String) {
|
||||
self.logs.push(log);
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
|
||||
pub enum TaskResult {
|
||||
Success(String),
|
||||
Error(String),
|
||||
Split(Vec<ToDo>),
|
||||
}
|
||||
|
||||
impl TaskResult {
|
||||
pub fn result_type(&self) -> TaskResultType {
|
||||
match self {
|
||||
TaskResult::Success(_) => TaskResultType::Success,
|
||||
TaskResult::Error(_) => TaskResultType::Error,
|
||||
TaskResult::Split(_) => TaskResultType::Split,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn test_task_new() {
|
||||
let todo_id = Uuid::new_v4();
|
||||
let task = Task::new(todo_id, "Test task description".to_string());
|
||||
assert!(!task.id.is_nil());
|
||||
assert_eq!(task.todo_id, todo_id);
|
||||
assert_eq!(task.description, "Test task description");
|
||||
assert_eq!(task.status, TaskStatus::Pending);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_task_start() {
|
||||
let mut task = Task::default();
|
||||
assert!(task.start().is_ok());
|
||||
assert_eq!(task.status, TaskStatus::Running);
|
||||
assert!(task.started_at.is_some());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_task_complete() {
|
||||
let mut task = Task::default();
|
||||
task.status = TaskStatus::Running;
|
||||
let result = TaskResult::Success("output".to_string());
|
||||
assert!(task.complete(result.clone()).is_ok());
|
||||
assert_eq!(task.status, TaskStatus::Completed);
|
||||
assert!(task.completed_at.is_some());
|
||||
assert_eq!(task.result, Some(result));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_task_duration_ms() {
|
||||
let mut task = Task::default();
|
||||
task.start().expect("task should start");
|
||||
std::thread::sleep(std::time::Duration::from_millis(10));
|
||||
task.complete(TaskResult::Success("done".to_string()))
|
||||
.expect("task should complete");
|
||||
|
||||
let duration = task.duration_ms();
|
||||
assert!(duration.is_some());
|
||||
assert!(duration.unwrap() >= 10);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_task_retry() {
|
||||
let mut task = Task::default();
|
||||
task.status = TaskStatus::Failed;
|
||||
task.result = Some(TaskResult::Error("failed".to_string()));
|
||||
assert!(task.retry().is_ok());
|
||||
assert_eq!(task.status, TaskStatus::Pending);
|
||||
assert!(task.result.is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_task_default_values() {
|
||||
let task = Task::default();
|
||||
assert_eq!(task.status, TaskStatus::Pending);
|
||||
assert!(task.result.is_none());
|
||||
assert!(task.logs.is_empty());
|
||||
assert!(task.started_at.is_none());
|
||||
assert!(task.completed_at.is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_task_serialization_roundtrip() {
|
||||
let original = Task::new(Uuid::new_v4(), "Test task".to_string());
|
||||
let serialized = serde_json::to_string(&original).unwrap();
|
||||
let deserialized: Task = serde_json::from_str(&serialized).unwrap();
|
||||
|
||||
assert_eq!(original.id, deserialized.id);
|
||||
assert_eq!(original.todo_id, deserialized.todo_id);
|
||||
assert_eq!(original.description, deserialized.description);
|
||||
assert_eq!(original.status, deserialized.status);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_task_fail() {
|
||||
let mut task = Task::default();
|
||||
task.status = TaskStatus::Running;
|
||||
assert!(task.fail("error message".to_string()).is_ok());
|
||||
assert_eq!(task.status, TaskStatus::Failed);
|
||||
assert!(task.result.is_some());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_task_invalid_transition() {
|
||||
let mut task = Task::default();
|
||||
// Cannot complete a pending task directly
|
||||
let result = TaskResult::Success("done".to_string());
|
||||
assert!(task.complete(result).is_err());
|
||||
}
|
||||
}
|
||||
Loading…
Reference in a new issue