use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use std::collections::HashMap;
use std::path::PathBuf;
use std::time::Duration;
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub enum SessionStatus {
InProgress,
Completed,
Failed,
Interrupted,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub enum WorkflowType {
Standard,
Iterative,
StructuredWithOutputs,
MapReduce,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct SessionState {
pub session_id: String,
pub status: SessionStatus,
pub started_at: DateTime<Utc>,
pub ended_at: Option<DateTime<Utc>>,
pub iterations_completed: usize,
pub files_changed: usize,
pub errors: Vec<String>,
pub working_directory: PathBuf,
pub worktree_name: Option<String>,
pub workflow_started_at: Option<DateTime<Utc>>,
pub current_iteration_started_at: Option<DateTime<Utc>>,
pub current_iteration_number: Option<u32>,
pub iteration_timings: Vec<(u32, Duration)>,
pub command_timings: Vec<(String, Duration)>,
pub workflow_state: Option<WorkflowState>,
pub execution_environment: Option<ExecutionEnvironment>,
pub last_checkpoint: Option<DateTime<Utc>>,
pub workflow_hash: Option<String>,
pub workflow_type: Option<WorkflowType>,
pub execution_context: Option<ExecutionContext>,
pub checkpoint_version: u32,
pub last_validated_at: Option<DateTime<Utc>>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct WorkflowState {
pub current_iteration: usize,
pub current_step: usize,
pub completed_steps: Vec<StepResult>,
pub workflow_path: PathBuf,
pub input_args: Vec<String>,
pub map_patterns: Vec<String>,
pub using_worktree: bool,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct StepResult {
pub step_index: usize,
pub command: String,
pub success: bool,
pub output: Option<String>,
pub duration: Duration,
pub error: Option<String>,
pub started_at: DateTime<Utc>,
pub completed_at: DateTime<Utc>,
pub exit_code: Option<i32>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ExecutionContext {
pub variables: HashMap<String, String>,
pub step_outputs: HashMap<usize, String>,
pub environment: HashMap<String, String>,
}
impl Default for ExecutionContext {
fn default() -> Self {
Self::new()
}
}
impl ExecutionContext {
pub fn new() -> Self {
Self {
variables: HashMap::new(),
step_outputs: HashMap::new(),
environment: HashMap::new(),
}
}
pub fn restore_from_state(workflow_state: &WorkflowState) -> Self {
let mut context = Self::new();
for step in &workflow_state.completed_steps {
if let Some(ref output) = step.output {
context.step_outputs.insert(step.step_index, output.clone());
context
.variables
.insert(format!("step_{}_output", step.step_index), output.clone());
}
}
context
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ExecutionEnvironment {
pub working_directory: PathBuf,
pub worktree_name: Option<String>,
pub environment_vars: HashMap<String, String>,
pub command_args: Vec<String>,
}
impl SessionState {
pub fn new(session_id: String, working_directory: PathBuf) -> Self {
Self {
session_id,
status: SessionStatus::InProgress,
started_at: Utc::now(),
ended_at: None,
iterations_completed: 0,
files_changed: 0,
errors: Vec::new(),
working_directory,
worktree_name: None,
workflow_started_at: None,
current_iteration_started_at: None,
current_iteration_number: None,
iteration_timings: Vec::new(),
command_timings: Vec::new(),
workflow_state: None,
execution_environment: None,
last_checkpoint: None,
workflow_hash: None,
workflow_type: None,
execution_context: None,
checkpoint_version: 1,
last_validated_at: None,
}
}
pub fn complete(&mut self) {
self.status = SessionStatus::Completed;
self.ended_at = Some(Utc::now());
}
pub fn fail(&mut self, error: String) {
self.status = SessionStatus::Failed;
self.ended_at = Some(Utc::now());
self.errors.push(error);
}
pub fn interrupt(&mut self) {
self.status = SessionStatus::Interrupted;
self.ended_at = Some(Utc::now());
}
pub fn add_files_changed(&mut self, count: usize) {
self.files_changed += count;
}
pub fn increment_iteration(&mut self) {
self.iterations_completed += 1;
}
pub fn duration(&self) -> Option<chrono::Duration> {
self.ended_at
.map(|end| end.signed_duration_since(self.started_at))
}
pub fn is_resumable(&self) -> bool {
if matches!(self.status, SessionStatus::Completed) {
return false;
}
self.workflow_state.is_some()
}
pub fn update_workflow_state(&mut self, state: WorkflowState) {
self.workflow_state = Some(state);
self.last_checkpoint = Some(Utc::now());
}
pub fn get_resume_info(&self) -> Option<String> {
self.workflow_state.as_ref().map(|ws| {
format!(
"Step {}/{} in iteration {}",
ws.current_step + 1,
ws.completed_steps.len() + 1,
ws.current_iteration + 1
)
})
}
}