use std::path::{Path, PathBuf};
use tokio::process::Command;
use tracing::{debug, info, warn};
use crate::ai_command_runner::{AiCommandRunner, OutputLine};
use crate::config::OrchestratorConfig;
use crate::error::{OrchestratorError, Result};
use crate::vcs::git::commands as git_commands;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum RejectionReviewVerdict {
Confirm,
Resume,
Block,
}
fn rejected_file_path(workspace_path: &Path, change_id: &str) -> PathBuf {
workspace_path
.join("openspec")
.join("changes")
.join(change_id)
.join("REJECTED.md")
}
#[allow(dead_code)]
pub fn has_rejection_proposal(workspace_path: &Path, change_id: &str) -> bool {
rejected_file_path(workspace_path, change_id).is_file()
}
fn resolve_rejection_review_command(template: &str, prompt: &str, change_id: &str) -> String {
let command = OrchestratorConfig::expand_change_id(template, change_id);
OrchestratorConfig::expand_prompt(&command, prompt)
}
#[allow(dead_code)]
fn rejection_review_prompt(change_id: &str) -> String {
rejection_review_prompt_with_skill(
crate::config::defaults::DEFAULT_REJECTING_SKILL,
None,
change_id,
)
}
fn rejection_review_prompt_with_skill(
rejecting_skill: &str,
workspace_path: Option<&Path>,
change_id: &str,
) -> String {
format!(
"{}\n\nRejecting review id:{change_id}\n\nchange_id: {change_id}\nproposal_path: openspec/changes/{change_id}/proposal.md\ntasks_path: {tasks_path}\nrejected_path: openspec/changes/{change_id}/REJECTED.md",
crate::agent::prompt::skill_prelude(rejecting_skill),
tasks_path = crate::agent::prompt::selected_tasks_path(workspace_path, change_id)
)
}
fn parse_rejection_review_output(output: &str) -> Option<RejectionReviewVerdict> {
let mut in_code_block = false;
for line in output.lines() {
let trimmed = line.trim();
if trimmed.starts_with("```") {
in_code_block = !in_code_block;
continue;
}
if in_code_block {
continue;
}
if trimmed == "REJECTION_REVIEW: CONFIRM" {
return Some(RejectionReviewVerdict::Confirm);
}
if trimmed == "REJECTION_REVIEW: RESUME" {
return Some(RejectionReviewVerdict::Resume);
}
if trimmed == "REJECTION_REVIEW: BLOCK" {
return Some(RejectionReviewVerdict::Block);
}
}
None
}
pub async fn run_rejection_review(
change_id: &str,
workspace_path: &Path,
config: &OrchestratorConfig,
ai_runner: &AiCommandRunner,
) -> Result<RejectionReviewVerdict> {
let command_template = config.get_acceptance_command()?;
let prompt = rejection_review_prompt_with_skill(
config.get_rejecting_skill(),
Some(workspace_path),
change_id,
);
let command = resolve_rejection_review_command(command_template, &prompt, change_id);
info!(
change_id = %change_id,
workspace = %workspace_path.display(),
"Starting dedicated rejecting review"
);
let (mut child, mut output_rx) = ai_runner
.execute_streaming_with_retry(
&command,
Some(workspace_path),
Some("rejecting"),
Some(change_id),
)
.await?;
let mut stdout = String::new();
while let Some(line) = output_rx.recv().await {
match line {
OutputLine::Stdout(s) => {
stdout.push_str(&s);
stdout.push('\n');
}
OutputLine::Stderr(_) => {}
}
}
let status = child.wait().await.map_err(|e| {
OrchestratorError::AgentCommand(format!(
"Failed to wait for rejecting review command for change '{}': {}",
change_id, e
))
})?;
if !status.success() {
return Err(OrchestratorError::AgentCommand(format!(
"Rejecting review command failed with exit code {:?} for change '{}'",
status.code(),
change_id
)));
}
parse_rejection_review_output(&stdout).ok_or_else(|| {
OrchestratorError::AgentCommand(format!(
"Rejecting review output missing required marker for '{}' (expected exactly one of REJECTION_REVIEW: CONFIRM|RESUME|BLOCK)",
change_id
))
})
}
async fn clear_rejected_proposal_marker(change_id: &str, workspace_path: &Path) -> Result<PathBuf> {
let rejected_path = rejected_file_path(workspace_path, change_id);
if rejected_path.exists() {
tokio::fs::remove_file(&rejected_path).await?;
}
Ok(rejected_path)
}
async fn resolve_recovery_tasks_path(
change_id: &str,
workspace_path: &Path,
) -> Result<crate::task_file::TaskFile> {
let workspace_path = workspace_path.to_path_buf();
let owned_change_id = change_id.to_string();
let resolved = tokio::task::spawn_blocking(move || {
crate::task_file::resolve_mutation(&owned_change_id, &workspace_path)
})
.await
.map_err(|e| {
OrchestratorError::AgentCommand(format!(
"Failed to resolve rejecting recovery task artifact for '{}': {}",
change_id, e
))
})??;
let Some(resolved) = resolved else {
return Err(OrchestratorError::AgentCommand(format!(
"Failed to resolve the canonical task artifact for rejecting recovery '{}'. Expected {} or {} in the active or archived change entry.",
change_id,
crate::task_file::MARKDOWN_FILE_NAME,
crate::task_file::JSON_FILE_NAME
)));
};
debug!(
change_id = %change_id,
tasks_path = %resolved.file.display(),
location = resolved.kind.log_label(),
"Resolved rejecting recovery task artifact"
);
Ok(resolved.file)
}
async fn append_recovery_task(change_id: &str, workspace_path: &Path) -> Result<PathBuf> {
let tasks_file = resolve_recovery_tasks_path(change_id, workspace_path).await?;
let owned_change_id = change_id.to_string();
let written = tasks_file.clone();
tokio::task::spawn_blocking(move || {
crate::task_parser::append_rejecting_recovery_task(&written, &owned_change_id)
})
.await
.map_err(|e| {
OrchestratorError::AgentCommand(format!(
"Failed to update the task artifact while recording rejecting recovery for '{}' at '{}': {}",
change_id,
tasks_file.display(),
e
))
})??;
Ok(tasks_file.path)
}
pub async fn handle_resume_apply_from_rejecting(
change_id: &str,
workspace_path: &Path,
) -> Result<()> {
let rejected_path = clear_rejected_proposal_marker(change_id, workspace_path).await?;
let tasks_path = append_recovery_task(change_id, workspace_path).await?;
info!(
change_id = %change_id,
rejected_path = %rejected_path.display(),
tasks_path = %tasks_path.display(),
"Resumed apply from rejecting review"
);
Ok(())
}
pub async fn handle_blocked_from_rejecting(change_id: &str, workspace_path: &Path) -> Result<()> {
let rejected_path = clear_rejected_proposal_marker(change_id, workspace_path).await?;
let tasks_path = append_recovery_task(change_id, workspace_path).await?;
info!(
change_id = %change_id,
rejected_path = %rejected_path.display(),
tasks_path = %tasks_path.display(),
"Rejecting review returned BLOCK; workspace remains blocked"
);
Ok(())
}
fn rejected_markdown(change_id: &str, reason: &str) -> String {
format!(
"# REJECTED\n\n- change_id: {}\n- reason: {}\n",
change_id, reason
)
}
fn extract_rejected_reason(content: &str) -> Option<String> {
content
.lines()
.map(str::trim)
.find_map(|line| line.strip_prefix("- reason:").map(str::trim))
.filter(|reason| !reason.is_empty())
.map(ToString::to_string)
}
async fn resolve_rejection_reason(
workspace_path: &Path,
change_id: &str,
fallback_reason: &str,
) -> String {
let rejected_path = rejected_file_path(workspace_path, change_id);
match tokio::fs::read_to_string(&rejected_path).await {
Ok(content) => {
if let Some(reason) = extract_rejected_reason(&content) {
info!(
change_id = %change_id,
rejected_path = %rejected_path.display(),
"Using reason extracted from existing apply-generated REJECTED.md proposal"
);
reason
} else {
fallback_reason.to_string()
}
}
Err(_) => fallback_reason.to_string(),
}
}
async fn cleanup_worktree(repo_root: &Path, worktree_path: &Path) {
let worktree_path_str = worktree_path.to_string_lossy();
match git_commands::worktree_remove_with_options(
repo_root,
&worktree_path_str,
git_commands::WorktreeRemoveOptions::default(),
)
.await
{
Ok(()) => {
info!(
worktree = %worktree_path.display(),
repo_root = %repo_root.display(),
"Removed rejected worktree"
);
}
Err(e) => {
warn!(
error = %e,
worktree = %worktree_path.display(),
repo_root = %repo_root.display(),
"Failed to remove rejected worktree (may already be removed)"
);
}
}
}
pub async fn execute_rejection_flow(
change_id: &str,
reason: &str,
workspace_path: &Path,
base_branch: &str,
repo_root: &Path,
) -> Result<()> {
info!(
change_id = %change_id,
workspace = %workspace_path.display(),
repo_root = %repo_root.display(),
base_branch = %base_branch,
"Starting rejection flow"
);
let effective_reason = resolve_rejection_reason(workspace_path, change_id, reason).await;
git_commands::checkout(repo_root, base_branch)
.await
.map_err(OrchestratorError::from_vcs_error)?;
let rejected_path = rejected_file_path(repo_root, change_id);
let rejected_parent = rejected_path.parent().ok_or_else(|| {
OrchestratorError::AgentCommand(format!(
"Invalid REJECTED.md path for change '{}'",
change_id
))
})?;
tokio::fs::create_dir_all(rejected_parent).await?;
tokio::fs::write(
&rejected_path,
rejected_markdown(change_id, &effective_reason),
)
.await?;
let relative_rejected_path = format!("openspec/changes/{}/REJECTED.md", change_id);
let add_output = Command::new("git")
.args(["add", &relative_rejected_path])
.current_dir(repo_root)
.output()
.await?;
if !add_output.status.success() {
return Err(OrchestratorError::AgentCommand(format!(
"git add failed for '{}': {}",
relative_rejected_path,
String::from_utf8_lossy(&add_output.stderr).trim()
)));
}
let staged_paths_output = Command::new("git")
.args(["diff", "--cached", "--name-only"])
.current_dir(repo_root)
.output()
.await?;
if !staged_paths_output.status.success() {
return Err(OrchestratorError::AgentCommand(format!(
"git diff --cached --name-only failed for rejection '{}': {}",
change_id,
String::from_utf8_lossy(&staged_paths_output.stderr).trim()
)));
}
let staged_paths_stdout = String::from_utf8_lossy(&staged_paths_output.stdout).into_owned();
let staged_paths = staged_paths_stdout
.lines()
.map(str::trim)
.filter(|line| !line.is_empty())
.map(ToString::to_string)
.collect::<Vec<_>>();
if staged_paths != vec![relative_rejected_path.clone()] {
return Err(OrchestratorError::AgentCommand(format!(
"rejection flow staged unexpected files for '{}': {:?}",
change_id, staged_paths
)));
}
let commit_message = format!("reject(openspec): {}", change_id);
let commit_output = Command::new("git")
.args([
"commit",
"-m",
&commit_message,
"--",
&relative_rejected_path,
])
.current_dir(repo_root)
.output()
.await?;
if !commit_output.status.success() {
return Err(OrchestratorError::AgentCommand(format!(
"git commit failed for rejection '{}': {}",
change_id,
String::from_utf8_lossy(&commit_output.stderr).trim()
)));
}
debug!(change_id = %change_id, "Committed REJECTED.md on base branch");
cleanup_worktree(repo_root, workspace_path).await;
info!(change_id = %change_id, "Rejection flow completed");
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
use std::fs;
async fn init_git_repo(path: &Path) {
let status = Command::new("git")
.args(["init", "-b", "main"])
.current_dir(path)
.status()
.await
.expect("git init failed");
assert!(status.success());
let status = Command::new("git")
.args(["config", "user.email", "test@example.com"])
.current_dir(path)
.status()
.await
.expect("git config email failed");
assert!(status.success());
let status = Command::new("git")
.args(["config", "user.name", "Test User"])
.current_dir(path)
.status()
.await
.expect("git config name failed");
assert!(status.success());
fs::write(path.join("README.md"), "# test\n").expect("write readme");
let status = Command::new("git")
.args(["add", "."])
.current_dir(path)
.status()
.await
.expect("git add failed");
assert!(status.success());
let status = Command::new("git")
.args(["commit", "-m", "init"])
.current_dir(path)
.status()
.await
.expect("git commit failed");
assert!(status.success());
}
#[test]
fn test_rejection_review_prompt_uses_custom_skill_prelude() {
let prompt = rejection_review_prompt_with_skill("team-rejecting", None, "change-a");
assert!(prompt.contains("$team-rejecting"));
assert!(prompt.contains("load skills: team-rejecting"));
assert!(!prompt.contains("$cflx-rejecting"));
assert!(prompt.contains("Rejecting review id:change-a"));
}
#[test]
fn test_rejected_markdown_contains_reason() {
let content = rejected_markdown("change-a", "spec mismatch");
assert!(content.contains("change_id: change-a"));
assert!(content.contains("reason: spec mismatch"));
}
#[test]
fn test_extract_rejected_reason_parses_reason_line() {
let content = "# REJECTED\n\n- change_id: change-a\n- reason: apply blocked handoff\n";
let reason = extract_rejected_reason(content);
assert_eq!(reason.as_deref(), Some("apply blocked handoff"));
}
#[test]
fn test_extract_rejected_reason_returns_none_without_reason_line() {
let content = "# REJECTED\n\n- change_id: change-a\n";
let reason = extract_rejected_reason(content);
assert!(reason.is_none());
}
#[test]
fn test_rejected_file_path_layout() {
let path = rejected_file_path(Path::new("/tmp/ws"), "change-a");
assert!(path.ends_with("openspec/changes/change-a/REJECTED.md"));
}
#[test]
fn test_has_rejection_proposal_detects_marker_file() {
let temp_dir = tempfile::tempdir().unwrap();
let change_dir = temp_dir.path().join("openspec/changes/change-a");
std::fs::create_dir_all(&change_dir).unwrap();
assert!(
!has_rejection_proposal(temp_dir.path(), "change-a"),
"proposal should be absent before REJECTED.md exists"
);
std::fs::write(change_dir.join("REJECTED.md"), "# REJECTED").unwrap();
assert!(
has_rejection_proposal(temp_dir.path(), "change-a"),
"proposal should be detected after REJECTED.md is created"
);
}
#[test]
fn test_append_recovery_task_section_avoids_deleted_rejected_marker_reference() {
let temp_dir = tempfile::tempdir().expect("temp dir");
let tasks_file = crate::task_file::TaskFile::in_entry(
temp_dir.path(),
crate::task_file::TaskFileFormat::Markdown,
);
fs::write(
&tasks_file.path,
"## Implementation Tasks\n\n- [ ] keep going\n",
)
.expect("seed tasks");
crate::task_parser::append_rejecting_recovery_task(&tasks_file, "change-a")
.expect("recovery task is recorded");
let updated = fs::read_to_string(&tasks_file.path).expect("read updated tasks");
assert!(
updated.contains("## Rejecting Recovery Tasks"),
"recovery section should be appended"
);
assert!(
updated.contains("do not recreate REJECTED.md"),
"recovery task should explicitly avoid relying on deleted marker"
);
assert!(
!updated.contains("Investigate blocker in openspec/changes/change-a/REJECTED.md"),
"recovery task must not reference deleted worktree-local REJECTED.md"
);
}
#[tokio::test]
async fn test_resolve_recovery_tasks_path_prefers_active_tasks() {
let temp_dir = tempfile::tempdir().expect("temp dir");
let workspace = temp_dir.path();
let change_id = "change-a";
let active_dir = workspace.join("openspec").join("changes").join(change_id);
fs::create_dir_all(&active_dir).expect("create active dir");
fs::write(active_dir.join("tasks.md"), "## Implementation Tasks\n").expect("write tasks");
let resolved = resolve_recovery_tasks_path(change_id, workspace)
.await
.expect("resolve path");
assert_eq!(resolved.path, active_dir.join("tasks.md"));
}
#[tokio::test]
async fn test_resolve_recovery_tasks_path_falls_back_to_archive_tasks() {
let temp_dir = tempfile::tempdir().expect("temp dir");
let workspace = temp_dir.path();
let change_id = "change-a";
let archive_dir = workspace
.join("openspec")
.join("changes")
.join("archive")
.join("2026-04-29-change-a");
fs::create_dir_all(&archive_dir).expect("create archive dir");
fs::write(archive_dir.join("tasks.md"), "## Implementation Tasks\n").expect("write tasks");
let resolved = resolve_recovery_tasks_path(change_id, workspace)
.await
.expect("resolve path");
assert_eq!(resolved.path, archive_dir.join("tasks.md"));
}
#[tokio::test]
async fn test_resolve_recovery_tasks_path_reports_explored_paths_when_missing() {
let temp_dir = tempfile::tempdir().expect("temp dir");
let workspace = temp_dir.path();
let change_id = "change-a";
let archive_candidate = workspace
.join("openspec")
.join("changes")
.join("archive")
.join("2026-04-29-change-a");
fs::create_dir_all(&archive_candidate).expect("create archive candidate");
let err = resolve_recovery_tasks_path(change_id, workspace)
.await
.expect_err("expected path resolution failure");
let message = err.to_string();
assert!(
message.contains("rejecting recovery 'change-a'"),
"{message}"
);
assert!(message.contains("tasks.md"), "{message}");
assert!(message.contains("tasks.json"), "{message}");
}
#[tokio::test]
async fn test_handle_resume_apply_from_rejecting_updates_archived_tasks_file() {
let temp_dir = tempfile::tempdir().expect("temp dir");
let workspace = temp_dir.path();
let change_id = "change-a";
let archive_dir = workspace
.join("openspec")
.join("changes")
.join("archive")
.join("2026-04-29-change-a");
fs::create_dir_all(&archive_dir).expect("create archive dir");
fs::write(
archive_dir.join("tasks.md"),
"## Implementation Tasks\n- [ ] keep\n",
)
.expect("write tasks");
handle_resume_apply_from_rejecting(change_id, workspace)
.await
.expect("resume should succeed with archived tasks path");
let updated = fs::read_to_string(archive_dir.join("tasks.md")).expect("read updated tasks");
assert!(updated.contains("## Rejecting Recovery Tasks"));
}
#[tokio::test]
async fn test_handle_blocked_from_rejecting_updates_archived_tasks_file() {
let temp_dir = tempfile::tempdir().expect("temp dir");
let workspace = temp_dir.path();
let change_id = "change-a";
let archive_dir = workspace
.join("openspec")
.join("changes")
.join("archive")
.join("2026-04-29-change-a");
fs::create_dir_all(&archive_dir).expect("create archive dir");
fs::write(
archive_dir.join("tasks.md"),
"## Implementation Tasks\n- [ ] keep\n",
)
.expect("write tasks");
handle_blocked_from_rejecting(change_id, workspace)
.await
.expect("block should succeed with archived tasks path");
let updated = fs::read_to_string(archive_dir.join("tasks.md")).expect("read updated tasks");
assert!(updated.contains("## Rejecting Recovery Tasks"));
}
#[tokio::test]
async fn test_execute_rejection_flow_creates_marker_commits_and_cleans_worktree() {
let temp_dir = tempfile::tempdir().expect("temp dir");
let repo_root = temp_dir.path();
init_git_repo(repo_root).await;
let change_id = "blocked-change";
let change_dir = repo_root.join("openspec").join("changes").join(change_id);
fs::create_dir_all(&change_dir).expect("create change dir");
fs::write(change_dir.join("proposal.md"), "# proposal\n").expect("write proposal");
fs::write(change_dir.join("tasks.md"), "- [ ] task\n").expect("write tasks");
let current_branch = git_commands::get_current_branch(repo_root)
.await
.expect("current branch")
.expect("branch name");
let worktree_parent = repo_root.join(".worktrees");
fs::create_dir_all(&worktree_parent).expect("create worktree parent");
let worktree_path = worktree_parent.join(change_id);
git_commands::worktree_add(
repo_root,
worktree_path.to_str().expect("worktree path"),
&format!("wt/{}", change_id),
¤t_branch,
)
.await
.expect("create worktree");
let result = execute_rejection_flow(
change_id,
"Implementation blocker detected",
&worktree_path,
¤t_branch,
repo_root,
)
.await;
assert!(result.is_ok(), "rejection flow should succeed: {result:?}");
let marker_path = change_dir.join("REJECTED.md");
assert!(marker_path.exists(), "REJECTED.md must be created");
let marker = fs::read_to_string(&marker_path).expect("read marker");
assert!(marker.contains("change_id: blocked-change"));
assert!(marker.contains("reason: Implementation blocker detected"));
let head_message = Command::new("git")
.args(["log", "-1", "--pretty=%s"])
.current_dir(repo_root)
.output()
.await
.expect("read commit message");
assert!(head_message.status.success());
let message = String::from_utf8_lossy(&head_message.stdout);
assert!(message
.trim()
.starts_with("reject(openspec): blocked-change"));
let committed_paths = Command::new("git")
.args(["show", "--name-only", "--pretty=format:", "HEAD"])
.current_dir(repo_root)
.output()
.await
.expect("read committed paths");
assert!(committed_paths.status.success());
let committed_paths_stdout = String::from_utf8_lossy(&committed_paths.stdout).into_owned();
let committed_paths = committed_paths_stdout
.lines()
.map(str::trim)
.filter(|line| !line.is_empty())
.map(ToString::to_string)
.collect::<Vec<_>>();
assert_eq!(
committed_paths,
vec!["openspec/changes/blocked-change/REJECTED.md".to_string()],
"rejection commit must contain only REJECTED.md"
);
let list = git_commands::list_worktrees(repo_root)
.await
.expect("list worktrees after cleanup");
assert!(
!list
.iter()
.any(|(path, _, _, _, _)| path == &worktree_path.to_string_lossy()),
"rejected worktree must be removed"
);
}
}