use anyhow::{anyhow, Context, Result};
use std::path::PathBuf;
use tokio::fs;
pub async fn run_resume_workflow(
session_id: Option<String>,
_force: bool,
from_checkpoint: Option<String>,
_path: Option<PathBuf>,
) -> Result<()> {
let session_id = if let Some(id) = session_id {
id
} else {
return Err(anyhow!(
"No session ID provided. Please specify a session ID or job ID to resume.\n\
Use 'prodigy sessions list' to see available sessions.\n\
Use 'prodigy resume-job list' to see MapReduce jobs."
));
};
let resume_result = try_unified_resume(&session_id, from_checkpoint).await;
match resume_result {
Ok(()) => Ok(()),
Err(e) => {
Err(anyhow!(
"Failed to resume {}: {}\n\n\
Troubleshooting:\n\
- Check if the session/job exists: 'prodigy sessions list' or 'prodigy resume-job list'\n\
- Ensure the worktree hasn't been cleaned up\n\
- For MapReduce jobs, try: 'prodigy resume-job {}'\n\
- For regular workflows, ensure checkpoint files exist",
session_id,
e,
session_id
))
}
}
}
enum SessionType {
Workflow,
MapReduce,
}
async fn check_session_type(id: &str) -> Result<SessionType> {
let storage =
crate::storage::GlobalStorage::new().context("Failed to create global storage")?;
let session_manager = crate::unified_session::SessionManager::new(storage)
.await
.context("Failed to create session manager")?;
let session_id = crate::unified_session::SessionId::from_string(id.to_string());
let session = session_manager
.load_session(&session_id)
.await
.context("Session not found in UnifiedSessionManager")?;
match session.session_type {
crate::unified_session::SessionType::MapReduce => Ok(SessionType::MapReduce),
crate::unified_session::SessionType::Workflow => Ok(SessionType::Workflow),
}
}
async fn try_unified_resume(id: &str, from_checkpoint: Option<String>) -> Result<()> {
let id_type = detect_id_type(id);
match id_type {
IdType::SessionId => {
match try_resume_regular_workflow(id, from_checkpoint.clone()).await {
Ok(()) => Ok(()),
Err(e) => {
try_resume_mapreduce_from_session(id).await.or(Err(e))
}
}
}
IdType::MapReduceJobId => {
try_resume_mapreduce_job(id).await
}
IdType::Ambiguous => {
match check_session_type(id).await {
Ok(SessionType::Workflow) => {
try_resume_regular_workflow(id, from_checkpoint.clone()).await
}
Ok(SessionType::MapReduce) => {
try_resume_mapreduce_job(id).await
}
Err(_) => {
match try_resume_regular_workflow(id, from_checkpoint.clone()).await {
Ok(()) => Ok(()),
Err(e) => {
let error_msg = e.to_string();
if error_msg.contains("already completed")
|| error_msg.contains("was cancelled")
{
return Err(e);
}
try_resume_mapreduce_job(id).await
}
}
}
}
}
}
}
enum IdType {
SessionId, MapReduceJobId, Ambiguous, }
async fn find_worktree_for_session(worktrees_dir: &PathBuf, session_id: &str) -> Result<PathBuf> {
if !worktrees_dir.exists() {
return Err(anyhow!(
"Worktrees directory does not exist: {}",
worktrees_dir.display()
));
}
let mut repo_dirs = fs::read_dir(worktrees_dir)
.await
.context("Failed to read worktrees directory")?;
while let Some(repo_entry) = repo_dirs.next_entry().await? {
if !repo_entry.path().is_dir() {
continue;
}
let potential_worktree = repo_entry.path().join(session_id);
if potential_worktree.exists() {
return Ok(potential_worktree);
}
}
Err(anyhow!(
"Worktree not found for session: {}\n\
Searched in: {}\n\
The worktree may have been cleaned up. You cannot resume this session.",
session_id,
worktrees_dir.display()
))
}
fn detect_id_type(id: &str) -> IdType {
if id.starts_with("session-mapreduce-") || id.starts_with("mapreduce-") {
IdType::MapReduceJobId
} else if id.starts_with("session-") {
IdType::SessionId
} else {
IdType::Ambiguous
}
}
async fn try_resume_regular_workflow(
session_id: &str,
from_checkpoint: Option<String>,
) -> Result<()> {
let prodigy_home = crate::storage::get_default_storage_dir()
.context("Failed to determine Prodigy storage directory")?;
let lock_manager = crate::cook::execution::ResumeLockManager::new(prodigy_home.clone())
.context("Failed to create resume lock manager")?;
let _lock = lock_manager
.acquire_lock(session_id)
.await
.context("Failed to acquire resume lock")?;
let storage =
crate::storage::GlobalStorage::new().context("Failed to create global storage")?;
let session_manager = crate::unified_session::SessionManager::new(storage)
.await
.context("Failed to create session manager")?;
let session_id_obj = crate::unified_session::SessionId::from_string(session_id.to_string());
let session_data = if let Ok(session) = session_manager.load_session(&session_id_obj).await {
use crate::unified_session::SessionStatus;
match session.status {
SessionStatus::Completed => {
return Err(anyhow!(
"Session {} has already completed and cannot be resumed.\n\
There is nothing to resume for this session.",
session_id
));
}
SessionStatus::Cancelled => {
return Err(anyhow!(
"Session {} was cancelled and cannot be resumed.",
session_id
));
}
_ => {
}
}
Some(session)
} else {
None
};
let checkpoint_dir = prodigy_home
.join("state")
.join(session_id)
.join("checkpoints");
if !checkpoint_dir.exists() {
let error_context = if let Some(ref session) = session_data {
if let Some(error) = &session.error {
format!("\n\nThe session failed with:\n{}", error)
} else {
String::new()
}
} else {
String::new()
};
return Err(anyhow!(
"Cannot resume session {}: No checkpoints found.\n\
This workflow failed before any checkpoints were created.{}\n\n\
Checkpoint directory does not exist: {}\n\n\
You cannot resume from this failure. Please fix the issue and run the workflow again.",
session_id,
error_context,
checkpoint_dir.display()
));
}
let has_checkpoints = std::fs::read_dir(&checkpoint_dir)
.ok()
.and_then(|entries| {
entries
.filter_map(Result::ok)
.any(|entry| entry.path().extension().and_then(|ext| ext.to_str()) == Some("json"))
.then_some(true)
})
.unwrap_or(false);
if !has_checkpoints {
let error_context = if let Some(ref session) = session_data {
if let Some(error) = &session.error {
format!("\n\nThe session failed with:\n{}", error)
} else {
String::new()
}
} else {
String::new()
};
return Err(anyhow!(
"Cannot resume session {}: No checkpoints found.\n\
This workflow failed before any checkpoints were created.{}\n\n\
Checkpoint directory exists but contains no checkpoint files: {}\n\n\
You cannot resume from this failure. Please fix the issue and run the workflow again.",
session_id,
error_context,
checkpoint_dir.display()
));
}
let checkpoint_file = if let Some(checkpoint_id) = &from_checkpoint {
let file = checkpoint_dir.join(format!("{}.checkpoint.json", checkpoint_id));
if !file.exists() {
return Err(anyhow!(
"Checkpoint not found: {}\nExpected at: {}",
checkpoint_id,
file.display()
));
}
file
} else {
let mut entries = fs::read_dir(&checkpoint_dir)
.await
.context("Failed to read checkpoint directory")?;
let mut latest_checkpoint: Option<(PathBuf, std::time::SystemTime)> = None;
while let Some(entry) = entries.next_entry().await? {
let path = entry.path();
if path.extension().and_then(|s| s.to_str()) == Some("json")
&& path
.file_name()
.and_then(|s| s.to_str())
.map(|s| s.ends_with(".checkpoint.json"))
.unwrap_or(false)
{
if let Ok(metadata) = entry.metadata().await {
if let Ok(modified) = metadata.modified() {
if latest_checkpoint.is_none()
|| modified > latest_checkpoint.as_ref().unwrap().1
{
latest_checkpoint = Some((path.clone(), modified));
}
}
}
}
}
latest_checkpoint
.ok_or_else(|| anyhow!("No checkpoint files found in {}", checkpoint_dir.display()))?
.0
};
let checkpoint_json = fs::read_to_string(&checkpoint_file)
.await
.with_context(|| {
format!(
"Failed to read checkpoint file: {}",
checkpoint_file.display()
)
})?;
let checkpoint: serde_json::Value =
serde_json::from_str(&checkpoint_json).context("Failed to parse checkpoint JSON")?;
let workflow_path = checkpoint
.get("workflow_path")
.and_then(|v| v.as_str())
.ok_or_else(|| {
anyhow!(
"Checkpoint does not contain workflow_path field.\n\n\
This checkpoint was created with an older version of Prodigy that didn't save\n\
the workflow file path. You can resume this session using:\n\n\
prodigy run <workflow-file>.yml --resume {}\n\n\
Where <workflow-file>.yml is the original workflow file you used.",
session_id
)
})?;
println!("🔄 Resuming session: {}", session_id);
println!("📄 Workflow: {}", workflow_path);
println!(
"📍 Checkpoint: {}",
checkpoint_file.file_name().unwrap().to_string_lossy()
);
let worktrees_dir = prodigy_home.join("worktrees");
let worktree_path = find_worktree_for_session(&worktrees_dir, session_id).await?;
println!();
println!("Note: Resuming from worktree: {}", worktree_path.display());
if from_checkpoint.is_some() {
println!(
" Using specific checkpoint: {}",
from_checkpoint.as_ref().unwrap()
);
} else {
println!(" Using latest checkpoint");
}
println!(" Project root: {}", worktree_path.display());
println!();
let workflow_pathbuf = PathBuf::from(workflow_path);
let cook_cmd = crate::cook::command::CookCommand {
playbook: workflow_pathbuf,
path: Some(worktree_path.clone()), max_iterations: 1,
map: vec![],
args: vec![],
fail_fast: false,
auto_accept: false,
resume: Some(session_id.to_string()), verbosity: 0,
quiet: false,
dry_run: false,
params: std::collections::HashMap::new(),
};
crate::cook::cook(cook_cmd).await
}
async fn try_resume_mapreduce_job(job_id: &str) -> Result<()> {
run_resume_job_command(job_id.to_string(), false, 0, None).await
}
async fn try_resume_mapreduce_from_session(session_id: &str) -> Result<()> {
let storage =
crate::storage::GlobalStorage::new().context("Failed to create global storage")?;
let session_manager = crate::unified_session::SessionManager::new(storage)
.await
.context("Failed to create session manager")?;
let session_id_obj = crate::unified_session::SessionId::from_string(session_id.to_string());
if let Ok(session) = session_manager.load_session(&session_id_obj).await {
use crate::unified_session::SessionStatus;
match session.status {
SessionStatus::Completed => {
return Err(anyhow!(
"Session {} has already completed and cannot be resumed.\n\
There is nothing to resume for this session.",
session_id
));
}
SessionStatus::Cancelled => {
return Err(anyhow!(
"Session {} was cancelled and cannot be resumed.",
session_id
));
}
_ => {
}
}
}
let prodigy_home = crate::storage::get_default_storage_dir()
.context("Failed to determine Prodigy storage directory")?;
let state_dir = prodigy_home.join("state");
if !state_dir.exists() {
return Err(anyhow!("No state directory found"));
}
let mut found_job_id: Option<String> = None;
if let Ok(entries) = fs::read_dir(&state_dir).await {
let mut entries = entries;
while let Ok(Some(repo_entry)) = entries.next_entry().await {
if !repo_entry.path().is_dir() {
continue;
}
let mapreduce_dir = repo_entry.path().join("mapreduce").join("jobs");
if !mapreduce_dir.exists() {
continue;
}
if let Ok(job_entries) = fs::read_dir(&mapreduce_dir).await {
let mut job_entries = job_entries;
while let Ok(Some(job_entry)) = job_entries.next_entry().await {
let job_name = job_entry.file_name();
let job_id = job_name.to_string_lossy();
if job_id.contains(session_id) {
found_job_id = Some(job_id.to_string());
break;
}
}
}
if found_job_id.is_some() {
break;
}
}
}
if let Some(job_id) = found_job_id {
println!("Found MapReduce job: {}", job_id);
try_resume_mapreduce_job(&job_id).await
} else {
Err(anyhow!(
"No MapReduce job found for session: {}",
session_id
))
}
}
pub async fn run_resume_job_command(
job_id: String,
_force: bool,
_max_retries: u32,
_path: Option<PathBuf>,
) -> Result<()> {
println!("🔄 Resuming MapReduce job: {}", job_id);
let prodigy_home = crate::storage::get_default_storage_dir()
.context("Failed to determine Prodigy storage directory")?;
let lock_manager = crate::cook::execution::ResumeLockManager::new(prodigy_home.clone())
.context("Failed to create resume lock manager")?;
let _lock = lock_manager
.acquire_lock(&job_id)
.await
.context("Failed to acquire resume lock")?;
let state_dir = prodigy_home.join("state");
if !state_dir.exists() {
return Err(anyhow!(
"Session not found: No state directory exists at: {}",
state_dir.display()
));
}
let mut job_path: Option<PathBuf> = None;
if let Ok(entries) = fs::read_dir(&state_dir).await {
let mut entries = entries;
while let Ok(Some(repo_entry)) = entries.next_entry().await {
if !repo_entry.path().is_dir() {
continue;
}
let potential_job_path = repo_entry
.path()
.join("mapreduce")
.join("jobs")
.join(&job_id);
if potential_job_path.exists() {
job_path = Some(potential_job_path);
break;
}
}
}
let job_dir = job_path.ok_or_else(|| {
anyhow!(
"MapReduce job not found: {}\n\
Searched in: {}",
job_id,
state_dir.display()
)
})?;
println!("📂 Found job at: {}", job_dir.display());
if let Ok(mut entries) = fs::read_dir(&job_dir).await {
println!("\n📋 Available checkpoints:");
while let Ok(Some(entry)) = entries.next_entry().await {
let name = entry.file_name();
if let Some(name_str) = name.to_str() {
if name_str.contains("checkpoint") {
println!(" - {}", name_str);
}
}
}
}
println!(
"\n🔍 Loading checkpoint and resuming execution for job: {}",
job_id
);
execute_mapreduce_resume(&job_id, _force, _max_retries, job_dir).await
}
async fn execute_mapreduce_resume(
job_id: &str,
force: bool,
max_retries: u32,
job_dir: PathBuf,
) -> Result<()> {
use crate::cook::execution::events::{EventLogger, JsonlEventWriter};
use crate::cook::execution::mapreduce_resume::{EnhancedResumeOptions, MapReduceResumeManager};
use crate::cook::execution::state::DefaultJobStateManager;
use crate::cook::orchestrator::ExecutionEnvironment;
use std::sync::Arc;
let state_dir = job_dir
.parent()
.and_then(|p| p.parent())
.and_then(|p| p.parent())
.ok_or_else(|| anyhow!("Invalid job directory structure"))?;
let repo_name = state_dir
.file_name()
.and_then(|n| n.to_str())
.ok_or_else(|| anyhow!("Could not determine repository name"))?;
let state_manager = Arc::new(DefaultJobStateManager::new(state_dir.to_path_buf()));
let checkpoint = state_manager
.checkpoint_manager
.load_checkpoint(job_id)
.await
.context("Failed to load checkpoint")?;
let working_dir = if let Some(parent_worktree) = &checkpoint.parent_worktree {
PathBuf::from(parent_worktree)
} else {
let prodigy_home = crate::storage::get_default_storage_dir()?;
let worktrees_dir = prodigy_home.join("worktrees").join(repo_name);
if !worktrees_dir.exists() {
return Err(anyhow!(
"No parent worktree found in checkpoint and worktrees directory does not exist: {}",
worktrees_dir.display()
));
}
let mut entries = fs::read_dir(&worktrees_dir)
.await
.context("Failed to read worktrees directory")?;
let mut newest_worktree: Option<PathBuf> = None;
let mut newest_time = std::time::SystemTime::UNIX_EPOCH;
while let Ok(Some(entry)) = entries.next_entry().await {
if entry.path().is_dir() {
if let Ok(metadata) = entry.metadata().await {
if let Ok(modified) = metadata.modified() {
if modified > newest_time {
newest_time = modified;
newest_worktree = Some(entry.path());
}
}
}
}
}
newest_worktree.ok_or_else(|| {
anyhow!(
"No worktrees found in: {}. The MapReduce job may have been cleaned up.",
worktrees_dir.display()
)
})?
};
println!("📂 Working directory: {}", working_dir.display());
println!("📊 Job has {} total items", checkpoint.total_items);
println!("✅ Completed: {}", checkpoint.successful_count);
println!("❌ Failed: {}", checkpoint.failed_count);
println!(
"⏳ Remaining: {}",
checkpoint.total_items - checkpoint.successful_count - checkpoint.failed_count
);
let events_dir = job_dir.join("events");
tokio::fs::create_dir_all(&events_dir)
.await
.context("Failed to create events directory")?;
let event_writer = Box::new(
JsonlEventWriter::new(events_dir.join("events.jsonl"))
.await
.context("Failed to create event writer")?,
);
let event_logger = Arc::new(EventLogger::new(vec![event_writer]));
let resume_manager = MapReduceResumeManager::new(
job_id.to_string(),
state_manager.clone(),
event_logger.clone(),
state_dir.to_path_buf(),
)
.await
.context("Failed to create resume manager")?;
let options = EnhancedResumeOptions {
force,
max_additional_retries: max_retries,
skip_validation: false,
from_checkpoint: None,
max_parallel: None,
force_recreation: false,
include_dlq_items: true,
validate_environment: true,
reset_failed_agents: false,
};
let env = ExecutionEnvironment {
working_dir: Arc::new(working_dir.clone()),
project_dir: Arc::new(working_dir.clone()),
worktree_name: None,
session_id: Arc::from(job_id),
};
println!("\n🚀 Starting resume execution...\n");
let result = resume_manager
.resume_job(job_id, options, &env)
.await
.context("Failed to resume MapReduce job")?;
display_resume_summary(&result)?;
Ok(())
}
fn display_resume_summary(
result: &crate::cook::execution::mapreduce_resume::EnhancedResumeResult,
) -> Result<()> {
use crate::cook::execution::mapreduce_resume::EnhancedResumeResult;
println!("\n");
println!("═══════════════════════════════════════════════");
println!(" MapReduce Resume Summary");
println!("═══════════════════════════════════════════════");
match result {
EnhancedResumeResult::FullWorkflowCompleted(full_result) => {
println!("\n✅ Workflow completed successfully!");
println!("\nMap Phase:");
println!(" • Total items: {}", full_result.map_result.total);
println!(" • Successful: {}", full_result.map_result.successful);
println!(" • Failed: {}", full_result.map_result.failed);
if let Some(reduce_result) = &full_result.reduce_result {
println!("\nReduce Phase:");
println!(" • Output: {}", reduce_result);
}
}
EnhancedResumeResult::MapOnlyCompleted(map_result) => {
println!("\n✅ Map phase completed!");
println!("\nResults:");
println!(" • Total items: {}", map_result.total);
println!(" • Successful: {}", map_result.successful);
println!(" • Failed: {}", map_result.failed);
println!("\n⚠️ Note: No reduce phase defined in workflow");
}
EnhancedResumeResult::PartialResume { phase, progress } => {
println!("\n⚠️ Partial resume (interrupted)");
println!("\nStatus:");
println!(" • Phase: {:?}", phase);
println!(" • Progress: {:.1}%", progress * 100.0);
println!("\n💡 Run 'prodigy resume-job <job_id>' again to continue");
}
EnhancedResumeResult::ReadyToExecute {
phase,
remaining_items,
state,
..
} => {
println!("\n⚠️ Resume prepared but not executed");
println!("\nStatus:");
println!(" • Phase: {:?}", phase);
println!(" • Remaining items: {}", remaining_items.len());
println!(" • Completed: {}", state.completed_agents.len());
println!("\n💡 Note: This indicates the resume manager prepared the state but did not execute");
println!(" This may occur if execution was not triggered properly");
}
}
println!("\n═══════════════════════════════════════════════\n");
Ok(())
}