diff --git a/Cargo.lock b/Cargo.lock index c6a5cd0..ad860bb 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1021,7 +1021,7 @@ checksum = "5d99f8c9a7727884afe522e9bd5edbfc91a3312b36a77b5fb8926e4c31a41801" [[package]] name = "tutti" -version = "0.2.0" +version = "0.2.2" dependencies = [ "chrono", "clap", diff --git a/Cargo.toml b/Cargo.toml index 2df6a4c..4f9297e 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "tutti" -version = "0.2.0" +version = "0.2.2" edition = "2024" # intentional: codebase uses Rust 2024 let-chain syntax description = "Multi-agent orchestration CLI — your agents, all together" license = "MIT" diff --git a/src/automation/mod.rs b/src/automation/mod.rs index 012ecfd..3b3048b 100644 --- a/src/automation/mod.rs +++ b/src/automation/mod.rs @@ -8,9 +8,10 @@ use crate::health::WaitFailureReason; use crate::permissions::evaluate_command_policy; use crate::session::TmuxSession; use crate::state::{ - AutomationRunRecord, ControlEvent, VerifyLastSummary, append_automation_run, - append_control_event, append_policy_decision, load_workflow_checkpoint, - save_verify_last_summary, save_workflow_checkpoint, save_workflow_output, + AutomationRunRecord, ControlEvent, VerifyLastSummary, WorkflowStepIntentRecord, + WorkflowStepOutcomeRecord, append_automation_run, append_control_event, append_policy_decision, + load_workflow_checkpoint, load_workflow_intent, save_verify_last_summary, + save_workflow_checkpoint, save_workflow_intent, save_workflow_output, }; use chrono::Utc; use serde::{Deserialize, Serialize}; @@ -519,6 +520,40 @@ fn step_is_control(step: &ResolvedStep) -> bool { ) } +fn step_type_name(step: &ResolvedStep) -> &'static str { + match step { + ResolvedStep::Prompt { .. } => "prompt", + ResolvedStep::Command { .. } => "command", + ResolvedStep::EnsureRunning { .. } => "ensure_running", + ResolvedStep::Workflow { .. } => "workflow", + ResolvedStep::Land { .. } => "land", + ResolvedStep::Review { .. } => "review", + } +} + +fn sanitize_step_key(input: &str) -> String { + input + .chars() + .map(|ch| { + if ch.is_ascii_alphanumeric() || matches!(ch, '.' | '_' | '-') { + ch + } else { + '_' + } + }) + .collect() +} + +fn step_file_key(step: &ResolvedStep, step_index: usize) -> String { + let base = match step { + ResolvedStep::Prompt { step_id, .. } | ResolvedStep::Command { step_id, .. } => { + step_id.as_deref().unwrap_or_else(|| step_type_name(step)) + } + _ => step_type_name(step), + }; + format!("{step_index:03}-{}", sanitize_step_key(base)) +} + #[derive(Debug, Clone)] pub struct PromptInjectedFile { source: PathBuf, @@ -575,6 +610,7 @@ pub struct ResumeContext { pub strict: bool, pub agent_scope: Option, pub started_at: chrono::DateTime, + pub failed_steps: Vec, pub completed_steps: HashSet, pub step_results: Vec, pub output_files: HashMap, @@ -622,6 +658,7 @@ pub fn load_resume_context(project_root: &Path, run_id: &str) -> Result WorkflowExecutor<'a> { let completed_steps = resume .map(|r| r.completed_steps.clone()) .unwrap_or_default(); + let mut attempted_steps = HashSet::new(); let explicit_dep_mode = workflow .steps .iter() @@ -684,12 +722,17 @@ impl<'a> WorkflowExecutor<'a> { let dag = execute_control_dag( self.config, self.project_root, + &run_id, + &workflow.name, &workflow.steps, &dep_graph, &completed_steps, )?; success = dag.success; failed_steps.extend(dag.failed_steps); + for result in &dag.step_results { + attempted_steps.insert(result.index); + } step_results.extend(dag.step_results); } else { for (idx, step) in workflow.steps.iter().enumerate() { @@ -697,6 +740,30 @@ impl<'a> WorkflowExecutor<'a> { if completed_steps.contains(&step_index) { continue; } + let prior_unsuccessful_intent = if resume.is_some() { + load_unsuccessful_step_intent(self.project_root, &run_id, step, step_index)? + } else { + None + }; + attempted_steps.insert(step_index); + if let Err(err) = + record_step_intent(self.project_root, &run_id, &workflow.name, step_index, step) + { + failed_steps.push(step_index); + success = false; + step_results.push(StepResult { + index: step_index, + step_type: step_type_name(step).to_string(), + status: StepStatus::Failed, + duration_ms: 0, + exit_code: None, + timed_out: false, + message: Some(format!("failed to write step intent: {err}")), + stdout: None, + stderr: None, + }); + break; + } match step { ResolvedStep::Prompt { step_id, @@ -1333,6 +1400,27 @@ impl<'a> WorkflowExecutor<'a> { .. } => { let started = std::time::Instant::now(); + if let Some(intent) = prior_unsuccessful_intent.as_ref() + && let Some(message) = should_skip_land_replay( + self.config, + self.project_root, + agent, + intent, + ) + { + step_results.push(StepResult { + index: step_index, + step_type: "land".to_string(), + status: StepStatus::Success, + duration_ms: started.elapsed().as_millis() as u64, + exit_code: Some(0), + timed_out: false, + message: Some(message), + stdout: None, + stderr: None, + }); + continue; + } if let Err(e) = ensure_agent_session_running(self.config, self.project_root, agent) { @@ -1426,6 +1514,29 @@ impl<'a> WorkflowExecutor<'a> { .. } => { let started = std::time::Instant::now(); + if let Some(intent) = prior_unsuccessful_intent.as_ref() + && let Some(packet) = review_packet_exists_since( + self.project_root, + agent, + intent.planned_at, + )? + { + step_results.push(StepResult { + index: step_index, + step_type: "review".to_string(), + status: StepStatus::Success, + duration_ms: started.elapsed().as_millis() as u64, + exit_code: Some(0), + timed_out: false, + message: Some(format!( + "resume guard: existing review packet detected at {}; skipping replay", + packet.display() + )), + stdout: None, + stderr: None, + }); + continue; + } if let Err(e) = ensure_agent_session_running(self.config, self.project_root, reviewer) { @@ -1527,6 +1638,26 @@ impl<'a> WorkflowExecutor<'a> { output_files, }; + for step_index in attempted_steps { + if let Some(step) = workflow.steps.get(step_index.saturating_sub(1)) + && let Some(step_result) = + result.step_results.iter().find(|s| s.index == step_index) + && let Err(err) = record_step_outcome( + self.project_root, + &result.run_id, + &workflow.name, + step_index, + step, + step_result, + ) + { + eprintln!( + "warn: failed to persist workflow step outcome for run '{}' step {}: {}", + result.run_id, step_index, err + ); + } + } + append_automation_run( self.project_root, &AutomationRunRecord { @@ -1786,6 +1917,360 @@ fn load_resume_outputs( Ok(outputs) } +fn step_intent_payload(step: &ResolvedStep) -> Value { + match step { + ResolvedStep::Prompt { + agent, + text, + session_name, + .. + } => json!({ + "agent": agent, + "session_name": session_name, + "text_chars": text.chars().count(), + }), + ResolvedStep::Command { + run, + cwd, + agent, + timeout_secs, + .. + } => json!({ + "run": run, + "cwd": cwd.display().to_string(), + "agent": agent, + "timeout_secs": timeout_secs, + }), + ResolvedStep::EnsureRunning { + agent, + session_name, + .. + } => json!({ + "agent": agent, + "session_name": session_name, + }), + ResolvedStep::Workflow { + workflow, + agent_override, + strict, + .. + } => json!({ + "workflow": workflow, + "agent_override": agent_override, + "strict": strict, + }), + ResolvedStep::Land { + agent, pr, force, .. + } => json!({ + "agent": agent, + "pr": pr, + "force": force, + }), + ResolvedStep::Review { + agent, reviewer, .. + } => json!({ + "agent": agent, + "reviewer": reviewer, + }), + } +} + +fn record_step_intent( + project_root: &Path, + run_id: &str, + workflow_name: &str, + step_index: usize, + step: &ResolvedStep, +) -> Result<()> { + let step_id = step_file_key(step, step_index); + let existing = load_workflow_intent(project_root, run_id, &step_id)?; + let attempt = existing + .as_ref() + .map(|record| record.attempt.saturating_add(1)) + .unwrap_or(1); + let planned_at = existing + .as_ref() + .map(|record| record.planned_at) + .unwrap_or_else(Utc::now); + let record = WorkflowStepIntentRecord { + run_id: run_id.to_string(), + workflow_name: workflow_name.to_string(), + step_index, + step_id: step_id.clone(), + step_type: step_type_name(step).to_string(), + planned_at, + intent: step_intent_payload(step), + attempt, + outcome: None, + }; + save_workflow_intent(project_root, run_id, &step_id, &record)?; + Ok(()) +} + +fn step_side_effects(step: &ResolvedStep, result: &StepResult) -> Option { + match step { + ResolvedStep::Land { + agent, pr, force, .. + } => Some(json!({ + "agent": agent, + "pr": pr, + "force": force, + "status": format!("{:?}", result.status).to_lowercase(), + })), + ResolvedStep::Review { + agent, reviewer, .. + } => Some(json!({ + "agent": agent, + "reviewer": reviewer, + "status": format!("{:?}", result.status).to_lowercase(), + })), + ResolvedStep::EnsureRunning { agent, .. } => Some(json!({ + "agent": agent, + "status": format!("{:?}", result.status).to_lowercase(), + })), + _ => None, + } +} + +fn record_step_outcome( + project_root: &Path, + run_id: &str, + workflow_name: &str, + step_index: usize, + step: &ResolvedStep, + result: &StepResult, +) -> Result<()> { + let step_id = step_file_key(step, step_index); + let mut record = + load_workflow_intent(project_root, run_id, &step_id)?.unwrap_or(WorkflowStepIntentRecord { + run_id: run_id.to_string(), + workflow_name: workflow_name.to_string(), + step_index, + step_id: step_id.clone(), + step_type: step_type_name(step).to_string(), + planned_at: Utc::now(), + intent: step_intent_payload(step), + attempt: 1, + outcome: None, + }); + record.workflow_name = workflow_name.to_string(); + record.step_index = step_index; + record.step_type = step_type_name(step).to_string(); + record.outcome = Some(WorkflowStepOutcomeRecord { + completed_at: Utc::now(), + status: format!("{:?}", result.status).to_lowercase(), + success: result.status == StepStatus::Success, + exit_code: result.exit_code, + timed_out: result.timed_out, + message: result.message.clone(), + side_effects: step_side_effects(step, result), + }); + save_workflow_intent(project_root, run_id, &step_id, &record)?; + Ok(()) +} + +fn load_unsuccessful_step_intent( + project_root: &Path, + run_id: &str, + step: &ResolvedStep, + step_index: usize, +) -> Result> { + let step_id = step_file_key(step, step_index); + let Some(record) = load_workflow_intent(project_root, run_id, &step_id)? else { + return Ok(None); + }; + if record.outcome.as_ref().is_some_and(|o| o.success) { + return Ok(None); + } + Ok(Some(record)) +} + +fn review_packet_exists_since( + project_root: &Path, + agent: &str, + since: chrono::DateTime, +) -> Result> { + let reviews_dir = project_root.join(".tutti").join("state").join("reviews"); + if !reviews_dir.exists() { + return Ok(None); + } + + let prefix = format!("{agent}-review-"); + let mut best: Option<(std::time::SystemTime, PathBuf)> = None; + for entry in std::fs::read_dir(&reviews_dir)? { + let entry = entry?; + let path = entry.path(); + if !path.is_file() { + continue; + } + let Some(name) = path.file_name().and_then(|f| f.to_str()) else { + continue; + }; + if !name.starts_with(&prefix) || !name.ends_with(".md") { + continue; + } + let Ok(meta) = entry.metadata() else { + continue; + }; + let Ok(modified) = meta.modified() else { + continue; + }; + let modified_utc: chrono::DateTime = modified.into(); + if modified_utc < since { + continue; + } + if best + .as_ref() + .is_none_or(|(best_modified, _)| modified > *best_modified) + { + best = Some((modified, path)); + } + } + Ok(best.map(|(_, path)| path)) +} + +fn agent_branch(config: &TuttiConfig, agent: &str) -> Option { + config + .agents + .iter() + .find(|a| a.name == agent) + .map(|a| a.resolved_branch()) +} + +fn is_branch_merged(project_root: &Path, branch: &str) -> Result { + let status = Command::new("git") + .args(["merge-base", "--is-ancestor", branch, "HEAD"]) + .current_dir(project_root) + .status()?; + Ok(status.success()) +} + +fn worktree_has_changes(project_root: &Path, agent: &str) -> Result { + let snapshot = crate::worktree::inspect_worktree(project_root, agent)?; + Ok(snapshot.exists && snapshot.dirty) +} + +fn should_skip_land_replay( + config: &TuttiConfig, + project_root: &Path, + agent: &str, + intent: &WorkflowStepIntentRecord, +) -> Option { + let branch = agent_branch(config, agent).unwrap_or_else(|| format!("tutti/{agent}")); + let merged = is_branch_merged(project_root, &branch).unwrap_or(false); + let dirty_worktree = worktree_has_changes(project_root, agent).unwrap_or(true); + if merged && !dirty_worktree { + Some(format!( + "resume guard: '{}' already merged (branch '{}') with no pending worktree changes after prior attempt at {}; skipping replay", + agent, + branch, + intent.planned_at.to_rfc3339() + )) + } else { + None + } +} + +fn land_replay_diagnostics( + config: &TuttiConfig, + project_root: &Path, + agent: &str, +) -> (String, bool, bool) { + let branch = agent_branch(config, agent).unwrap_or_else(|| format!("tutti/{agent}")); + let merged = is_branch_merged(project_root, &branch).unwrap_or(false); + let dirty_worktree = worktree_has_changes(project_root, agent).unwrap_or(true); + (branch, merged, dirty_worktree) +} + +pub fn build_resume_compensator_plan( + config: &TuttiConfig, + project_root: &Path, + workflow: &ResolvedWorkflow, + resume: &ResumeContext, +) -> Result> { + let Some(step_index) = resume + .failed_steps + .iter() + .copied() + .find(|idx| !resume.completed_steps.contains(idx)) + else { + return Ok(Vec::new()); + }; + let Some(step) = workflow.steps.get(step_index.saturating_sub(1)) else { + return Ok(Vec::new()); + }; + let Some(intent) = + load_unsuccessful_step_intent(project_root, &resume.run_id, step, step_index)? + else { + return Ok(Vec::new()); + }; + + let mut plan = vec![format!( + "Compensator preflight: step {step_index} ({}) has prior intent without a successful outcome.", + step_type_name(step) + )]; + + match step { + ResolvedStep::Land { agent, .. } => { + let (branch, merged, dirty_worktree) = + land_replay_diagnostics(config, project_root, agent); + if merged && !dirty_worktree { + plan.push(format!( + "Detected branch '{branch}' already merged into HEAD with no pending worktree changes; replay will be skipped by idempotency guard." + )); + } else { + plan.push(format!( + "Branch '{branch}' is not safely idempotent yet (merged={merged}, worktree_dirty={dirty_worktree}); replay will re-attempt `land`." + )); + } + } + ResolvedStep::Review { agent, .. } => { + if let Some(packet) = + review_packet_exists_since(project_root, agent, intent.planned_at)? + { + plan.push(format!( + "Detected existing review packet after prior attempt: {}. Replay will be skipped by idempotency guard.", + packet.display() + )); + } else { + plan.push( + "No new review packet found after prior attempt; replay will re-send review." + .to_string(), + ); + } + } + ResolvedStep::EnsureRunning { + agent, + session_name, + .. + } => { + if TmuxSession::session_exists(session_name) { + plan.push(format!( + "Agent '{agent}' is already running; replay remains a no-op." + )); + } else { + plan.push(format!( + "Agent '{agent}' is not running; replay will attempt to start it." + )); + } + } + ResolvedStep::Command { .. } => { + plan.push( + "Command steps are best-effort only; verify repository side effects before replay." + .to_string(), + ); + } + _ => { + plan.push( + "Replay will continue; review prior step output if side effects are possible." + .to_string(), + ); + } + } + + Ok(plan) +} + fn save_execution_checkpoint( project_root: &Path, options: &ExecuteOptions, @@ -1892,6 +2377,8 @@ fn build_normalized_dependencies(steps: &[ResolvedStep]) -> Vec> { fn execute_control_dag( config: &TuttiConfig, project_root: &Path, + run_id: &str, + workflow_name: &str, steps: &[ResolvedStep], dependencies: &[Vec], completed_steps: &HashSet, @@ -1926,6 +2413,8 @@ fn execute_control_dag( vec![execute_control_step( config, project_root, + run_id, + workflow_name, &steps[ready[0] - 1], ready[0], )?] @@ -1934,10 +2423,12 @@ fn execute_control_dag( for idx in &ready { let cfg = config.clone(); let root = project_root.to_path_buf(); + let run_id = run_id.to_string(); + let workflow_name = workflow_name.to_string(); let step = steps[*idx - 1].clone(); let step_idx = *idx; handles.push(std::thread::spawn(move || { - execute_control_step(&cfg, &root, &step, step_idx) + execute_control_step(&cfg, &root, &run_id, &workflow_name, &step, step_idx) })); } let mut out = Vec::with_capacity(ready.len()); @@ -1981,10 +2472,31 @@ fn execute_control_dag( fn execute_control_step( config: &TuttiConfig, project_root: &Path, + run_id: &str, + workflow_name: &str, step: &ResolvedStep, step_index: usize, ) -> Result { let started = std::time::Instant::now(); + let prior_unsuccessful_intent = + load_unsuccessful_step_intent(project_root, run_id, step, step_index)?; + if let Err(err) = record_step_intent(project_root, run_id, workflow_name, step_index, step) { + return Ok(ControlStepOutcome { + index: step_index, + result: StepResult { + index: step_index, + step_type: step_type_name(step).to_string(), + status: StepStatus::Failed, + duration_ms: started.elapsed().as_millis() as u64, + exit_code: None, + timed_out: false, + message: Some(format!("failed to write step intent: {err}")), + stdout: None, + stderr: None, + }, + hard_fail: true, + }); + } match step { ResolvedStep::EnsureRunning { agent, @@ -2065,6 +2577,29 @@ fn execute_control_step( fail_mode, .. } => { + if let Some(intent) = prior_unsuccessful_intent.as_ref() + && let Some(packet) = + review_packet_exists_since(project_root, agent, intent.planned_at)? + { + return Ok(ControlStepOutcome { + index: step_index, + result: StepResult { + index: step_index, + step_type: "review".to_string(), + status: StepStatus::Success, + duration_ms: started.elapsed().as_millis() as u64, + exit_code: Some(0), + timed_out: false, + message: Some(format!( + "resume guard: existing review packet detected at {}; skipping replay", + packet.display() + )), + stdout: None, + stderr: None, + }, + hard_fail: false, + }); + } let reviewer_session = TmuxSession::session_name(&config.workspace.name, reviewer); if !TmuxSession::session_exists(&reviewer_session) && let Err(e) = @@ -2118,6 +2653,25 @@ fn execute_control_step( fail_mode, .. } => { + if let Some(intent) = prior_unsuccessful_intent.as_ref() + && let Some(message) = should_skip_land_replay(config, project_root, agent, intent) + { + return Ok(ControlStepOutcome { + index: step_index, + result: StepResult { + index: step_index, + step_type: "land".to_string(), + status: StepStatus::Success, + duration_ms: started.elapsed().as_millis() as u64, + exit_code: Some(0), + timed_out: false, + message: Some(message), + stdout: None, + stderr: None, + }, + hard_fail: false, + }); + } let agent_session = TmuxSession::session_name(&config.workspace.name, agent); if !TmuxSession::session_exists(&agent_session) && let Err(e) = run_tt_subcommand(project_root, &["up".to_string(), agent.clone()]) @@ -2611,7 +3165,8 @@ mod tests { use super::*; use crate::config::{AgentConfig, DefaultsConfig, PermissionsConfig, WorkspaceConfig}; use std::collections::HashMap; - use std::path::PathBuf; + use std::path::{Path, PathBuf}; + use std::process::Command; fn sample_config(workflow: WorkflowConfig, hooks: Vec) -> TuttiConfig { TuttiConfig { @@ -2647,6 +3202,44 @@ mod tests { } } + fn init_git_repo(dir: &Path) { + git_run(dir, &["init"]); + git_run(dir, &["config", "user.email", "tutti-tests@example.com"]); + git_run(dir, &["config", "user.name", "Tutti Tests"]); + std::fs::write(dir.join("README.md"), "ok\n").unwrap(); + git_run(dir, &["add", "README.md"]); + git_run(dir, &["commit", "-m", "init"]); + } + + fn git_run(dir: &Path, args: &[&str]) { + let output = Command::new("git") + .args(args) + .current_dir(dir) + .output() + .unwrap(); + assert!( + output.status.success(), + "git {} failed: {}", + args.join(" "), + String::from_utf8_lossy(&output.stderr) + ); + } + + fn git_stdout(dir: &Path, args: &[&str]) -> String { + let output = Command::new("git") + .args(args) + .current_dir(dir) + .output() + .unwrap(); + assert!( + output.status.success(), + "git {} failed: {}", + args.join(" "), + String::from_utf8_lossy(&output.stderr) + ); + String::from_utf8_lossy(&output.stdout).trim().to_string() + } + #[test] fn command_fail_open_continues() { let workflow = WorkflowConfig { @@ -3251,6 +3844,371 @@ mod tests { let _ = std::fs::remove_dir_all(&dir); } + #[test] + fn execute_persists_step_intent_and_outcome() { + let workflow = WorkflowConfig { + name: "verify".to_string(), + description: None, + schedule: None, + steps: vec![WorkflowStepConfig::Command { + id: None, + depends_on: vec![], + run: "echo ok".to_string(), + cwd: Some(WorkflowCommandCwd::Workspace), + subdir: None, + agent: None, + timeout_secs: Some(30), + fail_mode: Some(WorkflowFailMode::Closed), + output_json: None, + }], + }; + + let dir = std::env::temp_dir().join("tutti-test-step-intent-outcome"); + let _ = std::fs::remove_dir_all(&dir); + std::fs::create_dir_all(&dir).unwrap(); + crate::state::ensure_tutti_dir(&dir).unwrap(); + + let config = sample_config(workflow, vec![]); + let opts = ExecuteOptions { + strict: false, + force_open_commands: false, + command_policy: None, + retry_policy: None, + origin: ExecutionOrigin::Run, + hook_event: None, + hook_agent: None, + }; + let resolved = WorkflowResolver::new(&config, &dir) + .resolve("verify", None, &opts) + .unwrap(); + let result = WorkflowExecutor::new(&config, &dir) + .execute(&resolved, &opts, None, None, None) + .unwrap(); + + let intent = crate::state::load_workflow_intent(&dir, &result.run_id, "001-command") + .unwrap() + .unwrap(); + assert_eq!(intent.step_type, "command"); + assert!( + intent + .outcome + .as_ref() + .is_some_and(|outcome| outcome.success) + ); + + let _ = std::fs::remove_dir_all(&dir); + } + + #[test] + fn record_step_intent_preserves_original_planned_at() { + let dir = std::env::temp_dir().join("tutti-test-intent-preserve-planned-at"); + let _ = std::fs::remove_dir_all(&dir); + std::fs::create_dir_all(&dir).unwrap(); + crate::state::ensure_tutti_dir(&dir).unwrap(); + + let step = ResolvedStep::Command { + step_id: None, + depends_on: vec![], + run: "echo ok".to_string(), + cwd: dir.clone(), + agent: None, + timeout_secs: 30, + fail_mode: WorkflowFailMode::Closed, + output_json: None, + }; + + record_step_intent(&dir, "run1", "verify", 1, &step).unwrap(); + let first = crate::state::load_workflow_intent(&dir, "run1", "001-command") + .unwrap() + .unwrap(); + std::thread::sleep(std::time::Duration::from_millis(2)); + record_step_intent(&dir, "run1", "verify", 1, &step).unwrap(); + let second = crate::state::load_workflow_intent(&dir, "run1", "001-command") + .unwrap() + .unwrap(); + + assert_eq!(second.planned_at, first.planned_at); + assert_eq!(second.attempt, first.attempt + 1); + let _ = std::fs::remove_dir_all(&dir); + } + + #[test] + fn execute_control_step_reports_intent_write_failure_as_failed_result() { + let dir = std::env::temp_dir().join("tutti-test-dag-intent-failure"); + let _ = std::fs::remove_dir_all(&dir); + std::fs::create_dir_all(&dir).unwrap(); + // Force intent persistence to fail (`.tutti` is a file, not a directory). + std::fs::write(dir.join(".tutti"), "not-a-dir\n").unwrap(); + + let workflow = WorkflowConfig { + name: "autofix".to_string(), + description: None, + schedule: None, + steps: vec![WorkflowStepConfig::EnsureRunning { + depends_on: vec![], + agent: "backend".to_string(), + fail_mode: Some(WorkflowFailMode::Closed), + }], + }; + let config = sample_config(workflow, vec![]); + let step = ResolvedStep::EnsureRunning { + depends_on: vec![], + agent: "backend".to_string(), + session_name: "tutti-ws-backend".to_string(), + fail_mode: WorkflowFailMode::Closed, + }; + + let outcome = execute_control_step(&config, &dir, "run-fail", "autofix", &step, 1).unwrap(); + assert_eq!(outcome.result.status, StepStatus::Failed); + assert!(outcome.hard_fail); + assert!( + outcome + .result + .message + .as_deref() + .is_some_and(|m| m.contains("failed to write step intent")) + ); + + let _ = std::fs::remove_dir_all(&dir); + } + + #[test] + fn build_resume_compensator_plan_warns_for_command_step() { + let workflow = WorkflowConfig { + name: "verify".to_string(), + description: None, + schedule: None, + steps: vec![WorkflowStepConfig::Command { + id: None, + depends_on: vec![], + run: "exit 1".to_string(), + cwd: Some(WorkflowCommandCwd::Workspace), + subdir: None, + agent: None, + timeout_secs: Some(30), + fail_mode: Some(WorkflowFailMode::Closed), + output_json: None, + }], + }; + let dir = std::env::temp_dir().join("tutti-test-resume-plan-command"); + let _ = std::fs::remove_dir_all(&dir); + std::fs::create_dir_all(&dir).unwrap(); + crate::state::ensure_tutti_dir(&dir).unwrap(); + + let config = sample_config(workflow, vec![]); + let checkpoint = serde_json::json!({ + "run_id": "run-command", + "workflow_name": "verify", + "strict": false, + "origin": "run", + "agent_scope": null, + "started_at": Utc::now(), + "finished_at": Utc::now(), + "success": false, + "failed_steps": [1], + "step_results": [], + "output_files": {} + }); + crate::state::save_workflow_checkpoint(&dir, "run-command", &checkpoint).unwrap(); + + let intent = WorkflowStepIntentRecord { + run_id: "run-command".to_string(), + workflow_name: "verify".to_string(), + step_index: 1, + step_id: "001-command".to_string(), + step_type: "command".to_string(), + planned_at: Utc::now(), + intent: serde_json::json!({"run":"exit 1"}), + attempt: 1, + outcome: None, + }; + crate::state::save_workflow_intent(&dir, "run-command", "001-command", &intent).unwrap(); + + let resume = load_resume_context(&dir, "run-command").unwrap().unwrap(); + let opts = ExecuteOptions { + strict: false, + force_open_commands: false, + command_policy: None, + retry_policy: None, + origin: ExecutionOrigin::Run, + hook_event: None, + hook_agent: None, + }; + let resolved = WorkflowResolver::new(&config, &dir) + .resolve("verify", None, &opts) + .unwrap(); + let plan = build_resume_compensator_plan(&config, &dir, &resolved, &resume).unwrap(); + assert!( + plan.iter() + .any(|line| line.contains("best-effort") || line.contains("best effort")) + ); + + let _ = std::fs::remove_dir_all(&dir); + } + + #[test] + fn build_resume_compensator_plan_detects_already_merged_land_branch() { + let workflow = WorkflowConfig { + name: "autoland".to_string(), + description: None, + schedule: None, + steps: vec![WorkflowStepConfig::Land { + depends_on: vec![], + agent: "backend".to_string(), + pr: Some(false), + force: Some(false), + fail_mode: Some(WorkflowFailMode::Closed), + }], + }; + let dir = std::env::temp_dir().join("tutti-test-resume-plan-land"); + let _ = std::fs::remove_dir_all(&dir); + std::fs::create_dir_all(&dir).unwrap(); + init_git_repo(&dir); + let default_branch = git_stdout(&dir, &["rev-parse", "--abbrev-ref", "HEAD"]); + git_run(&dir, &["checkout", "-b", "tutti/backend"]); + std::fs::write(dir.join("feature.txt"), "hello\n").unwrap(); + git_run(&dir, &["add", "feature.txt"]); + git_run(&dir, &["commit", "-m", "feature"]); + git_run(&dir, &["checkout", &default_branch]); + git_run( + &dir, + &["merge", "--no-ff", "-m", "merge backend", "tutti/backend"], + ); + crate::state::ensure_tutti_dir(&dir).unwrap(); + + let config = sample_config(workflow, vec![]); + let checkpoint = serde_json::json!({ + "run_id": "run-land", + "workflow_name": "autoland", + "strict": false, + "origin": "run", + "agent_scope": "backend", + "started_at": Utc::now(), + "finished_at": Utc::now(), + "success": false, + "failed_steps": [1], + "step_results": [], + "output_files": {} + }); + crate::state::save_workflow_checkpoint(&dir, "run-land", &checkpoint).unwrap(); + + let intent = WorkflowStepIntentRecord { + run_id: "run-land".to_string(), + workflow_name: "autoland".to_string(), + step_index: 1, + step_id: "001-land".to_string(), + step_type: "land".to_string(), + planned_at: Utc::now(), + intent: serde_json::json!({"agent":"backend","pr":false,"force":false}), + attempt: 1, + outcome: None, + }; + crate::state::save_workflow_intent(&dir, "run-land", "001-land", &intent).unwrap(); + + let resume = load_resume_context(&dir, "run-land").unwrap().unwrap(); + let opts = ExecuteOptions { + strict: false, + force_open_commands: false, + command_policy: None, + retry_policy: None, + origin: ExecutionOrigin::Run, + hook_event: None, + hook_agent: None, + }; + let resolved = WorkflowResolver::new(&config, &dir) + .resolve("autoland", None, &opts) + .unwrap(); + let plan = build_resume_compensator_plan(&config, &dir, &resolved, &resume).unwrap(); + assert!(plan.iter().any(|line| line.contains("already merged"))); + + let _ = std::fs::remove_dir_all(&dir); + } + + #[test] + fn build_resume_compensator_plan_land_fresh_branch_is_not_skipped() { + let workflow = WorkflowConfig { + name: "autoland".to_string(), + description: None, + schedule: None, + steps: vec![WorkflowStepConfig::Land { + depends_on: vec![], + agent: "backend".to_string(), + pr: Some(false), + force: Some(false), + fail_mode: Some(WorkflowFailMode::Closed), + }], + }; + let dir = std::env::temp_dir().join("tutti-test-resume-plan-land-fresh"); + let _ = std::fs::remove_dir_all(&dir); + std::fs::create_dir_all(&dir).unwrap(); + init_git_repo(&dir); + git_run(&dir, &["branch", "tutti/backend"]); + crate::state::ensure_tutti_dir(&dir).unwrap(); + let worktree_path = dir.join(".tutti").join("worktrees").join("backend"); + git_run( + &dir, + &[ + "worktree", + "add", + worktree_path.to_str().unwrap(), + "tutti/backend", + ], + ); + std::fs::write(worktree_path.join("dirty.txt"), "pending\n").unwrap(); + + let config = sample_config(workflow, vec![]); + let checkpoint = serde_json::json!({ + "run_id": "run-land-fresh", + "workflow_name": "autoland", + "strict": false, + "origin": "run", + "agent_scope": "backend", + "started_at": Utc::now(), + "finished_at": Utc::now(), + "success": false, + "failed_steps": [1], + "step_results": [], + "output_files": {} + }); + crate::state::save_workflow_checkpoint(&dir, "run-land-fresh", &checkpoint).unwrap(); + + let intent = WorkflowStepIntentRecord { + run_id: "run-land-fresh".to_string(), + workflow_name: "autoland".to_string(), + step_index: 1, + step_id: "001-land".to_string(), + step_type: "land".to_string(), + planned_at: Utc::now(), + intent: serde_json::json!({"agent":"backend","pr":false,"force":false}), + attempt: 1, + outcome: None, + }; + crate::state::save_workflow_intent(&dir, "run-land-fresh", "001-land", &intent).unwrap(); + + let resume = load_resume_context(&dir, "run-land-fresh") + .unwrap() + .unwrap(); + let opts = ExecuteOptions { + strict: false, + force_open_commands: false, + command_policy: None, + retry_policy: None, + origin: ExecutionOrigin::Run, + hook_event: None, + hook_agent: None, + }; + let resolved = WorkflowResolver::new(&config, &dir) + .resolve("autoland", None, &opts) + .unwrap(); + let plan = build_resume_compensator_plan(&config, &dir, &resolved, &resume).unwrap(); + assert!( + plan.iter() + .any(|line| line.contains("not safely idempotent")) + ); + + let _ = std::fs::remove_dir_all(&dir); + } + #[test] fn resolver_applies_command_subdir() { let dir = std::env::temp_dir().join("tutti-test-command-subdir"); diff --git a/src/cli/run.rs b/src/cli/run.rs index 757d2d2..940fa4b 100644 --- a/src/cli/run.rs +++ b/src/cli/run.rs @@ -1,7 +1,7 @@ use crate::automation::{ ExecuteOptions, ExecutionOrigin, ExecutionResult, ResolvedStep, ResolvedWorkflow, StepStatus, - WorkflowResolver, execute_workflow_with_hooks, load_resume_context, - retry_policy_from_resilience, + WorkflowResolver, build_resume_compensator_plan, execute_workflow_with_hooks, + load_resume_context, retry_policy_from_resilience, }; use crate::config::{GlobalConfig, TuttiConfig}; use crate::error::{Result, TuttiError}; @@ -84,6 +84,13 @@ pub fn run( ))); } + if let Some(ctx) = resume_context.as_ref() { + let plan = build_resume_compensator_plan(&config, project_root, &resolved, ctx)?; + if !plan.is_empty() { + print_resume_plan(&ctx.run_id, &plan); + } + } + if dry_run { if json { println!( @@ -120,6 +127,13 @@ pub fn run( Ok(()) } +fn print_resume_plan(run_id: &str, plan: &[String]) { + eprintln!("Resume compensator plan for run '{run_id}':"); + for line in plan { + eprintln!(" - {line}"); + } +} + #[derive(Debug, Serialize)] struct WorkflowListItem { name: String, diff --git a/src/state/mod.rs b/src/state/mod.rs index c00454d..bcbff32 100644 --- a/src/state/mod.rs +++ b/src/state/mod.rs @@ -41,6 +41,35 @@ pub struct VerifyLastSummary { pub agent_scope: Option, } +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct WorkflowStepIntentRecord { + pub run_id: String, + pub workflow_name: String, + pub step_index: usize, + pub step_id: String, + pub step_type: String, + pub planned_at: DateTime, + pub intent: Value, + #[serde(default)] + pub attempt: u32, + #[serde(default)] + pub outcome: Option, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct WorkflowStepOutcomeRecord { + pub completed_at: DateTime, + pub status: String, + pub success: bool, + #[serde(default)] + pub exit_code: Option, + pub timed_out: bool, + #[serde(default)] + pub message: Option, + #[serde(default)] + pub side_effects: Option, +} + #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] #[serde(rename_all = "snake_case")] pub enum ActivityState { @@ -113,6 +142,7 @@ pub fn ensure_tutti_dir(project_root: &Path) -> Result { "state/runtime-settings", "state/health", "state/workflow-checkpoints", + "state/workflow-intents", "state/workflow-outputs", "worktrees", "handoffs", @@ -316,6 +346,43 @@ pub fn load_workflow_checkpoint( Ok(Some(value)) } +pub fn save_workflow_intent( + project_root: &Path, + run_id: &str, + step_id: &str, + record: &WorkflowStepIntentRecord, +) -> Result { + let dir = project_root + .join(".tutti") + .join("state") + .join("workflow-intents") + .join(run_id); + std::fs::create_dir_all(&dir)?; + let path = dir.join(format!("{step_id}.json")); + let body = serde_json::to_string_pretty(record)?; + std::fs::write(&path, body)?; + Ok(path) +} + +pub fn load_workflow_intent( + project_root: &Path, + run_id: &str, + step_id: &str, +) -> Result> { + let path = project_root + .join(".tutti") + .join("state") + .join("workflow-intents") + .join(run_id) + .join(format!("{step_id}.json")); + if !path.exists() { + return Ok(None); + } + let body = std::fs::read_to_string(path)?; + let record = serde_json::from_str(&body).map_err(|e| TuttiError::State(e.to_string()))?; + Ok(Some(record)) +} + pub fn append_control_event(project_root: &Path, event: &ControlEvent) -> Result<()> { let state_dir = project_root.join(".tutti").join("state"); std::fs::create_dir_all(&state_dir)?; @@ -518,6 +585,7 @@ mod tests { assert!(dir.join(".tutti/state").exists()); assert!(dir.join(".tutti/state/runtime-settings").exists()); assert!(dir.join(".tutti/state/workflow-checkpoints").exists()); + assert!(dir.join(".tutti/state/workflow-intents").exists()); assert!(dir.join(".tutti/worktrees").exists()); assert!(dir.join(".tutti/handoffs").exists()); assert!(dir.join(".tutti/logs").exists()); @@ -665,6 +733,43 @@ mod tests { std::fs::remove_dir_all(&dir).unwrap(); } + #[test] + fn workflow_intent_round_trip() { + let dir = + std::env::temp_dir().join(format!("tutti-test-workflow-intent-{}", std::process::id())); + ensure_tutti_dir(&dir).unwrap(); + + let record = WorkflowStepIntentRecord { + run_id: "run123".to_string(), + workflow_name: "verify".to_string(), + step_index: 1, + step_id: "step-001".to_string(), + step_type: "command".to_string(), + planned_at: Utc::now(), + intent: serde_json::json!({"run":"echo ok"}), + attempt: 1, + outcome: Some(WorkflowStepOutcomeRecord { + completed_at: Utc::now(), + status: "success".to_string(), + success: true, + exit_code: Some(0), + timed_out: false, + message: None, + side_effects: None, + }), + }; + let path = save_workflow_intent(&dir, "run123", "step-001", &record).unwrap(); + assert!(path.exists()); + let loaded = load_workflow_intent(&dir, "run123", "step-001") + .unwrap() + .unwrap(); + assert_eq!(loaded.workflow_name, "verify"); + assert_eq!(loaded.step_type, "command"); + assert!(loaded.outcome.as_ref().is_some_and(|o| o.success)); + + std::fs::remove_dir_all(&dir).unwrap(); + } + #[test] fn control_events_append_and_load() { let dir =