use crate::config::WorkflowConfig;
use crate::cook::orchestrator::environment::OrchestratorEnv;
use crate::cook::orchestrator::pure;
use crate::core::orchestration::{ExecutionMode, ExecutionPlan, Phase};
use std::collections::HashMap;
use std::path::PathBuf;
use std::sync::Arc;
use stillwater::{fail, from_async, pure as stillwater_pure, BoxedEffect, EffectExt};
pub type OrchEffect<T> = BoxedEffect<T, anyhow::Error, OrchestratorEnv>;
#[derive(Debug, Clone)]
pub struct WorkflowSession {
pub session_id: String,
pub worktree_path: Option<PathBuf>,
pub config: WorkflowConfig,
}
#[derive(Debug, Clone)]
pub struct StepContext {
pub session_id: String,
pub step_index: usize,
pub working_dir: PathBuf,
pub previous_results: Vec<StepResult>,
}
#[derive(Debug, Clone)]
pub struct StepResult {
pub step_index: usize,
pub success: bool,
pub output: String,
pub error: Option<String>,
}
#[derive(Debug, Clone)]
pub struct WorkflowResult {
pub session_id: String,
pub step_results: Vec<StepResult>,
pub success: bool,
}
#[derive(Debug, Clone)]
pub struct ExecutionEnvironment {
pub working_dir: Arc<PathBuf>,
pub project_dir: Arc<PathBuf>,
pub worktree_name: Option<Arc<str>>,
pub session_id: Arc<str>,
pub variables: HashMap<String, String>,
}
pub fn setup_environment_effect(
plan: ExecutionPlan,
project_path: Arc<PathBuf>,
dry_run: bool,
) -> OrchEffect<ExecutionEnvironment> {
from_async(move |env: &OrchestratorEnv| {
let plan = plan.clone();
let project_path = project_path.clone();
let session_manager = env.session_manager.clone();
async move {
let session_id = crate::unified_session::SessionId::new().to_string();
let session_id_arc: Arc<str> = Arc::from(session_id.as_str());
session_manager
.start_session(&session_id)
.await
.map_err(|e| anyhow::anyhow!("Failed to start session: {}", e))?;
let (working_dir, worktree_name) = if plan.requires_worktrees() && !dry_run {
(project_path.clone(), None)
} else {
(project_path.clone(), None)
};
Ok(ExecutionEnvironment {
working_dir,
project_dir: project_path,
worktree_name,
session_id: session_id_arc,
variables: HashMap::new(),
})
}
})
.boxed()
}
pub fn execute_plan_effect(
plan: ExecutionPlan,
exec_env: ExecutionEnvironment,
) -> OrchEffect<WorkflowResult> {
from_async(move |_env: &OrchestratorEnv| {
let plan = plan.clone();
let exec_env = exec_env.clone();
async move {
let step_results = match plan.mode {
ExecutionMode::MapReduce => {
vec![]
}
ExecutionMode::Standard | ExecutionMode::Iterative => {
execute_phases(&plan.phases, &exec_env).await?
}
ExecutionMode::DryRun => {
vec![StepResult {
step_index: 0,
success: true,
output: "Dry run completed".to_string(),
error: None,
}]
}
};
let success = step_results.iter().all(|r| r.success);
Ok(WorkflowResult {
session_id: exec_env.session_id.to_string(),
step_results,
success,
})
}
})
.boxed()
}
async fn execute_phases(
phases: &[Phase],
_exec_env: &ExecutionEnvironment,
) -> anyhow::Result<Vec<StepResult>> {
let mut results = Vec::new();
for (idx, phase) in phases.iter().enumerate() {
let result = StepResult {
step_index: idx,
success: true,
output: format!("Phase {} completed", phase),
error: None,
};
results.push(result);
}
Ok(results)
}
pub fn finalize_session_effect(
result: WorkflowResult,
exec_env: ExecutionEnvironment,
) -> OrchEffect<WorkflowResult> {
from_async(move |env: &OrchestratorEnv| {
let result = result.clone();
let exec_env = exec_env.clone();
let session_manager = env.session_manager.clone();
async move {
let _ = session_manager.complete_session().await;
if let Some(_worktree) = exec_env.worktree_name {
}
Ok(result)
}
})
.boxed()
}
pub fn run_workflow_effect(
plan: ExecutionPlan,
project_path: Arc<PathBuf>,
dry_run: bool,
) -> OrchEffect<WorkflowResult> {
let plan_for_exec = plan.clone();
let _plan_for_finalize = plan.clone();
setup_environment_effect(plan, project_path, dry_run)
.and_then_auto(move |exec_env| {
let env_for_finalize = exec_env.clone();
execute_plan_effect(plan_for_exec.clone(), exec_env)
.and_then_auto(move |result| finalize_session_effect(result, env_for_finalize))
})
.boxed()
}
pub fn validate_workflow(config: WorkflowConfig) -> OrchEffect<WorkflowConfig> {
use stillwater::Validation;
match pure::validate_workflow(&config) {
Validation::Success(_) => stillwater_pure(config).boxed(),
Validation::Failure(errors) => {
let error_msg = errors
.iter()
.map(|e| e.to_string())
.collect::<Vec<_>>()
.join(", ");
fail(anyhow::anyhow!("Workflow validation failed: {}", error_msg)).boxed()
}
}
}
pub fn setup_workflow(config: WorkflowConfig) -> OrchEffect<WorkflowSession> {
validate_workflow(config.clone())
.and_then_auto(|validated_config| {
from_async(move |env: &OrchestratorEnv| {
let session_manager = env.session_manager.clone();
async move {
let session_id = uuid::Uuid::new_v4().to_string();
session_manager
.start_session(&session_id)
.await
.map_err(|e| anyhow::anyhow!("Failed to start session: {}", e))?;
Ok::<WorkflowSession, anyhow::Error>(WorkflowSession {
session_id,
worktree_path: None,
config: validated_config,
})
}
})
})
.boxed()
}
pub fn execute_workflow(config: WorkflowConfig) -> OrchEffect<WorkflowResult> {
setup_workflow(config)
.and_then(|session| {
let result = WorkflowResult {
session_id: session.session_id,
step_results: vec![],
success: true,
};
stillwater_pure::<_, anyhow::Error, OrchestratorEnv>(result).boxed()
})
.boxed()
}
#[cfg(test)]
mod tests {
use super::*;
use crate::config::command::WorkflowCommand;
fn simple_workflow_config() -> WorkflowConfig {
WorkflowConfig {
name: Some("test".to_string()),
commands: vec![WorkflowCommand::Simple("echo test".to_string())],
env: None,
secrets: None,
env_files: None,
profiles: None,
merge: None,
}
}
fn empty_workflow_config() -> WorkflowConfig {
WorkflowConfig {
name: Some("empty".to_string()),
commands: vec![],
env: None,
secrets: None,
env_files: None,
profiles: None,
merge: None,
}
}
#[test]
fn test_workflow_session_creation() {
let config = simple_workflow_config();
let session = WorkflowSession {
session_id: "test-123".to_string(),
worktree_path: Some(PathBuf::from("/tmp/test")),
config,
};
assert_eq!(session.session_id, "test-123");
assert!(session.worktree_path.is_some());
}
#[test]
fn test_workflow_session_without_worktree() {
let config = simple_workflow_config();
let session = WorkflowSession {
session_id: "test-456".to_string(),
worktree_path: None,
config,
};
assert_eq!(session.session_id, "test-456");
assert!(session.worktree_path.is_none());
}
#[test]
fn test_step_context_creation() {
let context = StepContext {
session_id: "test-123".to_string(),
step_index: 0,
working_dir: PathBuf::from("/tmp"),
previous_results: vec![],
};
assert_eq!(context.session_id, "test-123");
assert_eq!(context.step_index, 0);
assert!(context.previous_results.is_empty());
}
#[test]
fn test_step_context_with_previous_results() {
let prev_result = StepResult {
step_index: 0,
success: true,
output: "previous output".to_string(),
error: None,
};
let context = StepContext {
session_id: "test-123".to_string(),
step_index: 1,
working_dir: PathBuf::from("/tmp"),
previous_results: vec![prev_result],
};
assert_eq!(context.step_index, 1);
assert_eq!(context.previous_results.len(), 1);
}
#[test]
fn test_step_result_success() {
let result = StepResult {
step_index: 0,
success: true,
output: "test output".to_string(),
error: None,
};
assert!(result.success);
assert!(result.error.is_none());
assert_eq!(result.output, "test output");
}
#[test]
fn test_step_result_failure() {
let result = StepResult {
step_index: 0,
success: false,
output: String::new(),
error: Some("test error".to_string()),
};
assert!(!result.success);
assert!(result.error.is_some());
assert_eq!(result.error.unwrap(), "test error");
}
#[test]
fn test_workflow_result_success() {
let result = WorkflowResult {
session_id: "test-789".to_string(),
step_results: vec![StepResult {
step_index: 0,
success: true,
output: "done".to_string(),
error: None,
}],
success: true,
};
assert!(result.success);
assert_eq!(result.step_results.len(), 1);
}
#[test]
fn test_workflow_result_failure() {
let result = WorkflowResult {
session_id: "test-999".to_string(),
step_results: vec![
StepResult {
step_index: 0,
success: true,
output: "step 1 ok".to_string(),
error: None,
},
StepResult {
step_index: 1,
success: false,
output: "".to_string(),
error: Some("step 2 failed".to_string()),
},
],
success: false,
};
assert!(!result.success);
assert_eq!(result.step_results.len(), 2);
assert!(!result.step_results[1].success);
}
#[test]
fn test_validate_workflow_pure_success() {
let config = simple_workflow_config();
let result = pure::validate_workflow(&config);
use stillwater::Validation;
assert!(matches!(result, Validation::Success(_)));
}
#[test]
fn test_validate_workflow_pure_failure() {
let config = empty_workflow_config();
let result = pure::validate_workflow(&config);
use stillwater::Validation;
assert!(matches!(result, Validation::Failure(_)));
}
#[test]
fn test_validate_workflow_effect_pure_transform() {
let config = simple_workflow_config();
let pure_result = pure::validate_workflow(&config);
use stillwater::Validation;
assert!(matches!(pure_result, Validation::Success(_)));
}
}