use crate::abstractions::git::GitOperations;
use crate::config::WorkflowConfig;
use crate::cook::execution::ClaudeExecutor;
use crate::cook::interaction::UserInteraction;
use crate::cook::orchestrator::core::{CookConfig, ExecutionEnvironment};
use crate::cook::session::{SessionManager, SessionState};
use crate::worktree::WorktreeManager;
use anyhow::{anyhow, Context, Result};
use sha2::{Digest, Sha256};
use std::sync::{Arc, Mutex};
#[derive(Clone)]
pub struct SessionOperations {
#[allow(dead_code)] session_manager: Arc<dyn SessionManager>,
claude_executor: Arc<dyn ClaudeExecutor>,
user_interaction: Arc<dyn UserInteraction>,
git_operations: Arc<dyn GitOperations>,
subprocess: crate::subprocess::SubprocessManager,
unified_session_manager: Arc<Mutex<Option<Arc<crate::unified_session::SessionManager>>>>,
}
impl SessionOperations {
pub fn new(
session_manager: Arc<dyn SessionManager>,
claude_executor: Arc<dyn ClaudeExecutor>,
user_interaction: Arc<dyn UserInteraction>,
git_operations: Arc<dyn GitOperations>,
subprocess: crate::subprocess::SubprocessManager,
) -> Self {
Self {
session_manager,
claude_executor,
user_interaction,
git_operations,
subprocess,
unified_session_manager: Arc::new(Mutex::new(None)),
}
}
pub async fn get_unified_session_manager(
&self,
) -> Result<Arc<crate::unified_session::SessionManager>> {
if let Some(m) = self
.unified_session_manager
.lock()
.map_err(|e| anyhow!("Lock: {}", e))?
.as_ref()
{
return Ok(Arc::clone(m));
}
let storage = crate::storage::GlobalStorage::new().context("Create storage")?;
let manager = Arc::new(
crate::unified_session::SessionManager::new(storage)
.await
.context("Create manager")?,
);
*self
.unified_session_manager
.lock()
.map_err(|e| anyhow!("Lock: {}", e))? = Some(Arc::clone(&manager));
Ok(manager)
}
pub async fn update_unified_session_status(&self, session_id: &str, success: bool) {
let Ok(manager) = self.get_unified_session_manager().await else {
return;
};
let id = crate::unified_session::SessionId::from_string(session_id.to_string());
if manager.load_session(&id).await.is_ok() {
let _ = manager.complete_session(&id, success).await;
}
}
pub async fn create_unified_session(&self, config: &CookConfig) -> Result<String> {
let manager = self.get_unified_session_manager().await?;
let wf_id = super::construction::generate_workflow_id();
let cfg = super::construction::build_session_config(
wf_id.clone(),
config.workflow.name.clone(),
super::construction::create_session_metadata(config.workflow.commands.len()),
);
let id = manager.create_session(cfg).await?;
manager.start_session(&id).await?;
log::info!("Created session: {} (workflow: {})", id, wf_id);
Ok(id.to_string())
}
pub fn generate_session_id(&self) -> String {
super::construction::generate_session_id()
}
pub fn calculate_workflow_hash(workflow: &WorkflowConfig) -> String {
let mut hasher = Sha256::new();
let serialized = serde_json::to_string(workflow).unwrap_or_default();
hasher.update(serialized);
format!("{:x}", hasher.finalize())
}
pub async fn check_prerequisites(&self) -> Result<()> {
let test_mode = std::env::var("PRODIGY_TEST_MODE").unwrap_or_default() == "true";
if test_mode {
return Ok(());
}
if !self.claude_executor.check_claude_cli().await? {
anyhow::bail!("Claude CLI is not available. Please install it first.");
}
if !self.git_operations.is_git_repo().await {
anyhow::bail!("Not in a git repository. Please run from a git repository.");
}
Ok(())
}
pub async fn check_prerequisites_with_config(&self, config: &CookConfig) -> Result<()> {
let test_mode = std::env::var("PRODIGY_TEST_MODE").unwrap_or_default() == "true";
if test_mode {
return Ok(());
}
if !self.claude_executor.check_claude_cli().await? {
anyhow::bail!("Claude CLI is not available. Please install it first.");
}
let is_temp_workflow = config
.command
.playbook
.to_str()
.map(|s| s.contains("/tmp/") || s.contains("/var/folders/") || s.contains("Temp"))
.unwrap_or(false);
if !is_temp_workflow && !self.git_operations.is_git_repo().await {
anyhow::bail!("Not in a git repository. Please run from a git repository.");
}
Ok(())
}
pub async fn restore_environment(
&self,
state: &SessionState,
config: &CookConfig,
) -> Result<ExecutionEnvironment> {
let mut working_dir = Arc::new(state.working_directory.clone());
let mut worktree_name: Option<Arc<str>> =
state.worktree_name.as_ref().map(|s| Arc::from(s.as_str()));
if let Some(ref name) = worktree_name {
let merge_config = config.workflow.merge.clone().or_else(|| {
config
.mapreduce_config
.as_ref()
.and_then(|m| m.merge.clone())
});
let workflow_env = config.workflow.env.clone().unwrap_or_default();
let worktree_manager = WorktreeManager::with_config(
config.project_path.to_path_buf(),
self.subprocess.clone(),
config.command.verbosity,
merge_config,
workflow_env,
)?;
let sessions = worktree_manager.list_sessions().await?;
if !sessions.iter().any(|s| s.name.as_str() == name.as_ref()) {
self.user_interaction
.display_warning(&format!("Worktree {} was deleted, recreating...", name));
let session = worktree_manager.create_session().await?;
working_dir = Arc::new(session.path.clone());
worktree_name = Some(Arc::from(session.name.as_ref()));
} else {
let sessions = worktree_manager.list_sessions().await?;
if let Some(session) = sessions.iter().find(|s| s.name.as_str() == name.as_ref()) {
working_dir = Arc::new(session.path.clone());
}
}
}
Ok(ExecutionEnvironment {
working_dir,
project_dir: Arc::clone(&config.project_path),
worktree_name,
session_id: Arc::from(state.session_id.as_str()),
})
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_calculate_workflow_hash_deterministic() {
use std::collections::HashMap;
let mut env1 = HashMap::new();
env1.insert("TEST".to_string(), "value".to_string());
let workflow = WorkflowConfig {
name: None,
commands: vec![],
env: Some(env1),
secrets: None,
env_files: None,
profiles: None,
merge: None,
};
let hash1 = SessionOperations::calculate_workflow_hash(&workflow);
let hash2 = SessionOperations::calculate_workflow_hash(&workflow);
assert_eq!(hash1, hash2, "Hash should be deterministic");
assert!(!hash1.is_empty(), "Hash should not be empty");
}
#[test]
fn test_calculate_workflow_hash_different_workflows() {
use std::collections::HashMap;
let mut env1 = HashMap::new();
env1.insert("TEST1".to_string(), "value1".to_string());
let mut env2 = HashMap::new();
env2.insert("TEST2".to_string(), "value2".to_string());
let workflow1 = WorkflowConfig {
name: None,
commands: vec![],
env: Some(env1),
secrets: None,
env_files: None,
profiles: None,
merge: None,
};
let workflow2 = WorkflowConfig {
name: None,
commands: vec![],
env: Some(env2),
secrets: None,
env_files: None,
profiles: None,
merge: None,
};
let hash1 = SessionOperations::calculate_workflow_hash(&workflow1);
let hash2 = SessionOperations::calculate_workflow_hash(&workflow2);
assert_ne!(
hash1, hash2,
"Different workflows should have different hashes"
);
}
}