use crate::agent::{AgentRunner, CleanupReviewDiagnostic, CleanupReviewFailureKind, OutputLine};
use crate::ai_command_runner::AiCommandRunner;
use crate::config::OrchestratorConfig;
use crate::error::{OrchestratorError, Result};
use crate::execution::apply as common_apply;
use crate::hooks::{HookContext, HookRunner, HookType};
use crate::parallel::output_bridge::ParallelApplyEventHandler;
use super::archive_state::delete_archive_state;
use super::events::ParallelEvent;
use crate::orchestration::build_acceptance_tail_findings;
use crate::stall::StallDetector;
use crate::vcs::git::commands as git_commands;
use crate::vcs::git::commands::has_uncommitted_changes;
use crate::vcs::git::GitWorkspaceManager;
use crate::vcs::VcsBackend;
use std::path::Path;
use std::sync::Arc;
use tokio::process::Command;
use tokio::sync::{mpsc, Mutex};
use tokio_util::sync::CancellationToken;
use tracing::{debug, error, info, warn};
async fn wait_for_streaming_child_with_cancel(
child: &mut crate::process_manager::StreamingChildHandle,
cancel_token: Option<&CancellationToken>,
operation: &str,
change_id: &str,
workspace_path: &Path,
attempt: Option<u32>,
) -> Result<std::process::ExitStatus> {
if let Some(token) = cancel_token {
tokio::select! {
_ = token.cancelled() => {
warn!(
operation = operation,
change_id = change_id,
workspace = %workspace_path.display(),
attempt = attempt,
"Cancellation observed while waiting for child status; terminating child"
);
let _ = child.terminate();
Err(OrchestratorError::cancelled(operation, change_id, workspace_path))
}
status = child.wait() => status.map_err(|e| {
let attempt_context = attempt
.map(|attempt| format!(" (attempt {})", attempt))
.unwrap_or_default();
OrchestratorError::AgentCommand(format!(
"Failed to wait for {} command for '{}' in workspace '{}'{}: {}",
operation,
change_id,
workspace_path.display(),
attempt_context,
e
))
}),
}
} else {
child.wait().await.map_err(|e| {
let attempt_context = attempt
.map(|attempt| format!(" (attempt {})", attempt))
.unwrap_or_default();
OrchestratorError::AgentCommand(format!(
"Failed to wait for {} command for '{}' in workspace '{}'{}: {}",
operation,
change_id,
workspace_path.display(),
attempt_context,
e
))
})
}
}
#[derive(Debug, Clone, Default)]
pub struct ParallelHookContext {
pub workspace_path: String,
pub group_index: Option<u32>,
#[allow(dead_code)] pub total_changes_in_group: usize,
pub total_changes: usize,
pub changes_processed: usize,
}
fn build_parallel_hook_context(
change_id: &str,
completed_tasks: u32,
total_tasks: u32,
apply_count: u32,
parallel_ctx: Option<&ParallelHookContext>,
) -> HookContext {
let (changes_processed, total_changes, remaining_changes) = match parallel_ctx {
Some(ctx) => (
ctx.changes_processed,
ctx.total_changes,
ctx.total_changes.saturating_sub(ctx.changes_processed),
),
None => (0, 0, 0),
};
let mut ctx = HookContext::new(changes_processed, total_changes, remaining_changes, false)
.with_change(change_id, completed_tasks, total_tasks)
.with_apply_count(apply_count);
if let Some(parallel_ctx) = parallel_ctx {
ctx = ctx.with_parallel_context(¶llel_ctx.workspace_path, parallel_ctx.group_index);
}
ctx
}
const MAX_CLEANUP_REVIEW_RETRIES: u32 = 2;
enum CleanupReviewAttempt {
Success,
Failure(CleanupReviewDiagnostic),
PermissionDenied(crate::permission::PermissionDenial),
}
async fn run_post_apply_cleanup_review(
change_id: &str,
workspace_path: &Path,
config: &OrchestratorConfig,
ai_runner: &AiCommandRunner,
cancel_token: Option<&CancellationToken>,
event_tx: Option<&mpsc::Sender<ParallelEvent>>,
) -> Result<()> {
let max_attempts = MAX_CLEANUP_REVIEW_RETRIES.saturating_add(1);
let mut latest: Option<CleanupReviewDiagnostic> = None;
for attempt in 1..=max_attempts {
if cancel_token.is_some_and(|token| token.is_cancelled()) {
return Err(OrchestratorError::cancelled(
"cleanup-review",
change_id,
workspace_path,
));
}
match run_cleanup_review_attempt(
change_id,
workspace_path,
config,
ai_runner,
cancel_token,
attempt,
max_attempts,
latest.as_ref(),
)
.await?
{
CleanupReviewAttempt::Success => {
info!(
change_id = %change_id,
workspace = %workspace_path.display(),
attempt = attempt,
"Post-apply cleanup review succeeded and worktree is clean"
);
return Ok(());
}
CleanupReviewAttempt::PermissionDenied(denial) => {
warn!(
change_id = %change_id,
category = denial.category.as_str(),
denied_target = %denial.denied_target,
"Cleanup-review blocked by permission/tool policy denial; entering non-terminal hold without a corrective attempt"
);
if let Some(tx) = event_tx {
let _ = tx
.send(ParallelEvent::ExecutionBlocked {
change_id: change_id.to_string(),
blocker: crate::events::StalledBlocker::permission_denial(
"cleanup-review",
&denial,
),
})
.await;
let _ = tx
.send(ParallelEvent::WorkspaceStatusUpdated {
change_id: change_id.to_string(),
workspace_name: workspace_path
.file_name()
.map(|name| name.to_string_lossy().to_string())
.unwrap_or_else(|| workspace_path.display().to_string()),
status: crate::vcs::WorkspaceStatus::Blocked,
})
.await;
}
return Err(OrchestratorError::PermissionStalled {
denied_path: denial.denied_target.clone(),
guidance: denial.format_guidance(),
});
}
CleanupReviewAttempt::Failure(diagnostic) => {
warn!(
change_id = %change_id,
workspace = %workspace_path.display(),
attempt = attempt,
max_attempts = max_attempts,
failure_kind = diagnostic.kind.label(),
marker_count = diagnostic.marker_count,
"Cleanup-review attempt did not produce a handoff-ready worktree"
);
if let Some(tx) = event_tx {
let _ = tx
.send(ParallelEvent::Log(
crate::events::LogEntry::warn(format!(
"Cleanup-review attempt {}/{} failed ({})",
attempt,
max_attempts,
diagnostic.kind.label()
))
.with_change_id(change_id)
.with_operation("cleanup-review"),
))
.await;
}
latest = Some(diagnostic);
}
}
}
let diagnostic = latest.expect("an exhausted cleanup-review loop recorded a failure");
Err(OrchestratorError::AgentCommand(format!(
"Cleanup-review failed on {} operation attempts for change '{}' in workspace '{}': {}",
max_attempts,
change_id,
workspace_path.display(),
format_cleanup_review_diagnostic(&diagnostic)
)))
}
fn format_cleanup_review_diagnostic(diagnostic: &CleanupReviewDiagnostic) -> String {
const MAX_TAIL_CHARS: usize = 400;
fn condense(tail: &str) -> String {
let single_line = tail.split_whitespace().collect::<Vec<_>>().join(" ");
match single_line.char_indices().nth(MAX_TAIL_CHARS) {
Some((idx, _)) => format!("{}...", &single_line[..idx]),
None => single_line,
}
}
let mut parts = vec![format!("failure_kind: {}", diagnostic.kind.label())];
if let Some(code) = diagnostic.exit_code {
parts.push(format!("exit_code: {}", code));
}
parts.push(format!(
"standalone_clean_marker_count: {}",
diagnostic.marker_count
));
if let Some(status) = diagnostic
.status_tail
.as_deref()
.filter(|tail| !tail.trim().is_empty())
{
parts.push(format!("status: {}", condense(status)));
}
if let Some(status_error) = diagnostic
.status_error
.as_deref()
.filter(|tail| !tail.trim().is_empty())
{
parts.push(format!("status_error: {}", condense(status_error)));
}
if let Some(stdout) = diagnostic
.stdout_tail
.as_deref()
.filter(|tail| !tail.trim().is_empty())
{
parts.push(format!("stdout: {}", condense(stdout)));
}
if let Some(stderr) = diagnostic
.stderr_tail
.as_deref()
.filter(|tail| !tail.trim().is_empty())
{
parts.push(format!("stderr: {}", condense(stderr)));
}
parts.join(" | ")
}
#[allow(clippy::too_many_arguments)]
async fn run_cleanup_review_attempt(
change_id: &str,
workspace_path: &Path,
config: &OrchestratorConfig,
ai_runner: &AiCommandRunner,
cancel_token: Option<&CancellationToken>,
attempt: u32,
max_attempts: u32,
previous: Option<&CleanupReviewDiagnostic>,
) -> Result<CleanupReviewAttempt> {
let user_template = config.get_acceptance_command()?;
let prompt = crate::agent::build_cleanup_review_prompt_with_skill(
config.get_cleanup_review_skill(),
Some(workspace_path),
change_id,
previous,
);
let command = OrchestratorConfig::expand_prompt(
&OrchestratorConfig::expand_change_id(user_template, change_id),
&prompt,
);
info!(
change_id = %change_id,
workspace = %workspace_path.display(),
attempt = attempt,
max_attempts = max_attempts,
"Starting post-apply cleanup review for dirty managed worktree"
);
let (mut child, mut output_rx) = ai_runner
.execute_streaming_with_retry(
&command,
Some(workspace_path),
Some("cleanup-review"),
Some(change_id),
)
.await?;
let mut output_collector = crate::history::OutputCollector::new();
let mut marker_scanner = crate::agent::CleanupMarkerScanner::new();
loop {
let line = if let Some(token) = cancel_token {
tokio::select! {
_ = token.cancelled() => {
warn!(
change_id = %change_id,
workspace = %workspace_path.display(),
attempt = attempt,
"Cancellation observed while streaming cleanup-review output; terminating child"
);
let _ = child.terminate();
output_rx.close();
while output_rx.recv().await.is_some() {}
return Err(OrchestratorError::cancelled(
"cleanup-review",
change_id,
workspace_path,
));
}
line = output_rx.recv() => line,
}
} else {
output_rx.recv().await
};
let Some(line) = line else { break };
match line {
crate::ai_command_runner::OutputLine::Stdout(s) => {
output_collector.add_stdout(&s);
marker_scanner.observe(&s);
}
crate::ai_command_runner::OutputLine::Stderr(s) => {
output_collector.add_stderr(&s);
}
}
}
let status = wait_for_streaming_child_with_cancel(
&mut child,
cancel_token,
"cleanup-review",
change_id,
workspace_path,
Some(attempt),
)
.await?;
let stdout_tail = output_collector.stdout_tail();
let stderr_tail = output_collector.stderr_tail();
let marker_count = marker_scanner.count();
if let Some(denial) = crate::permission::classify_permission_denial(&[
stdout_tail.as_deref(),
stderr_tail.as_deref(),
]) {
return Ok(CleanupReviewAttempt::PermissionDenied(denial));
}
if !status.success() {
let (status_tail, status_error) = split_status_inspection(workspace_path).await;
return Ok(CleanupReviewAttempt::Failure(CleanupReviewDiagnostic {
kind: CleanupReviewFailureKind::CommandFailed,
exit_code: status.code(),
stdout_tail,
stderr_tail,
marker_count,
status_tail,
status_error,
}));
}
if marker_count != 1 {
let (status_tail, status_error) = split_status_inspection(workspace_path).await;
let kind = if marker_count == 0 {
CleanupReviewFailureKind::MarkerMissing
} else {
CleanupReviewFailureKind::MarkerDuplicate
};
return Ok(CleanupReviewAttempt::Failure(CleanupReviewDiagnostic {
kind,
exit_code: status.code(),
stdout_tail,
stderr_tail,
marker_count,
status_tail,
status_error,
}));
}
match has_uncommitted_changes(workspace_path).await {
Ok((false, _)) => Ok(CleanupReviewAttempt::Success),
Ok((true, dirty_status)) => Ok(CleanupReviewAttempt::Failure(CleanupReviewDiagnostic {
kind: CleanupReviewFailureKind::DirtyRemains,
exit_code: status.code(),
stdout_tail,
stderr_tail,
marker_count,
status_tail: Some(dirty_status),
status_error: None,
})),
Err(e) => Ok(CleanupReviewAttempt::Failure(CleanupReviewDiagnostic {
kind: CleanupReviewFailureKind::StatusInspectionFailed,
exit_code: status.code(),
stdout_tail,
stderr_tail,
marker_count,
status_tail: None,
status_error: Some(format!("status inspection failed: {}", e)),
})),
}
}
async fn split_status_inspection(workspace_path: &Path) -> (Option<String>, Option<String>) {
match current_porcelain_status(workspace_path).await {
Ok(status) => (Some(status), None),
Err(e) => (None, Some(e.to_string())),
}
}
async fn current_porcelain_status(workspace_path: &Path) -> Result<String> {
match has_uncommitted_changes(workspace_path).await {
Ok((true, status)) => Ok(status),
Ok((false, _)) => Ok("clean".to_string()),
Err(e) => Err(OrchestratorError::GitCommand(format!(
"status inspection failed: {}",
e
))),
}
}
async fn mark_acceptance_context_injected(
agent: &AgentRunner,
change_id: &str,
acceptance_tail_injected: &Arc<Mutex<std::collections::HashMap<String, bool>>>,
) {
if agent.acceptance_context_was_injected(change_id) {
acceptance_tail_injected
.lock()
.await
.insert(change_id.to_string(), true);
}
}
async fn prepare_acceptance_context_for_apply(
agent: &mut AgentRunner,
change_id: &str,
acceptance_history: &Arc<Mutex<crate::history::AcceptanceHistory>>,
acceptance_tail_injected: &Arc<Mutex<std::collections::HashMap<String, bool>>>,
) {
agent.seed_acceptance_history(acceptance_history.lock().await.clone());
let already_injected = acceptance_tail_injected
.lock()
.await
.get(change_id)
.copied()
.unwrap_or(false);
if already_injected {
let _ = agent.get_acceptance_tail_context_for_apply(change_id);
}
}
#[allow(clippy::too_many_arguments)]
pub async fn execute_apply_in_workspace(
change_id: &str,
workspace_path: &Path,
_apply_cmd_template: &str,
config: &OrchestratorConfig,
event_tx: Option<mpsc::Sender<ParallelEvent>>,
vcs_backend: VcsBackend,
hooks: Option<&HookRunner>,
parallel_ctx: Option<&ParallelHookContext>,
cancel_token: Option<&CancellationToken>,
ai_runner: &AiCommandRunner,
repo_root: &Path,
_apply_history: &Arc<Mutex<crate::history::ApplyHistory>>,
acceptance_history: &Arc<Mutex<crate::history::AcceptanceHistory>>,
acceptance_tail_injected: &Arc<Mutex<std::collections::HashMap<String, bool>>>,
apply_budget: &common_apply::ApplyBudget,
) -> Result<(
String,
u32,
Option<crate::execution::apply::ApplyBlockedHandoff>,
Option<crate::execution::apply::ApplyRejectedHandoff>,
)> {
match git_commands::is_worktree(repo_root, workspace_path).await {
Ok(true) => {
info!(
"Workspace path validation passed: {} is a valid worktree",
workspace_path.display()
);
}
Ok(false) => {
let error_msg = format!(
"Parallel apply execution guard: workspace_path is NOT a worktree (executing in base repository is forbidden)
\
change_id: {}
\
workspace_path: {}
\
repo_root: {}",
change_id,
workspace_path.display(),
repo_root.display()
);
return Err(OrchestratorError::GitCommand(error_msg));
}
Err(e) => {
let error_msg = format!(
"Failed to validate worktree status for parallel apply
\
change_id: {}
\
workspace_path: {}
\
repo_root: {}
\
validation_error: {}",
change_id,
workspace_path.display(),
repo_root.display(),
e
);
return Err(OrchestratorError::GitCommand(error_msg));
}
}
let mut agent = AgentRunner::new(config.clone());
prepare_acceptance_context_for_apply(
&mut agent,
change_id,
acceptance_history,
acceptance_tail_injected,
)
.await;
let event_tx_for_permission = event_tx.clone();
let event_handler = ParallelApplyEventHandler::new(change_id.to_string(), event_tx);
let hook_ctx = match parallel_ctx {
Some(ctx) => common_apply::ApplyLoopHookContext::new(
ctx.changes_processed,
ctx.total_changes,
ctx.total_changes.saturating_sub(ctx.changes_processed),
workspace_path.to_string_lossy().to_string(),
ctx.group_index.unwrap_or(0) as usize,
),
None => common_apply::ApplyLoopHookContext::new(
0,
0,
0,
workspace_path.to_string_lossy().to_string(),
0,
),
};
let workspace_manager = GitWorkspaceManager::new(
workspace_path.parent().unwrap_or(repo_root).to_path_buf(),
repo_root.to_path_buf(),
1, config.clone(),
);
let apply_result = match common_apply::execute_apply_loop(
change_id,
workspace_path,
config,
&mut agent,
vcs_backend,
Some(&workspace_manager), hooks,
&hook_ctx,
&event_handler,
cancel_token,
ai_runner,
apply_budget,
|_line| async move {
},
)
.await
{
Ok(result) => result,
Err(crate::error::OrchestratorError::PermissionStalled {
denied_path,
guidance,
}) => {
warn!(
"Repeated unresolved permission/tool policy denial for {} in workspace {}: {}",
change_id,
workspace_path.display(),
denied_path
);
let denial = crate::permission::PermissionDenial {
category: crate::permission::PermissionDenialCategory::CommandPolicy,
denied_target: denied_path.clone(),
evidence: guidance.clone(),
};
if let Some(ref tx) = event_tx_for_permission {
let _ = tx
.send(ParallelEvent::ExecutionBlocked {
change_id: change_id.to_string(),
blocker: crate::events::StalledBlocker::permission_denial("apply", &denial),
})
.await;
let _ = tx
.send(ParallelEvent::WorkspaceStatusUpdated {
change_id: change_id.to_string(),
workspace_name: workspace_path
.file_name()
.map(|name| name.to_string_lossy().to_string())
.unwrap_or_else(|| workspace_path.display().to_string()),
status: crate::vcs::WorkspaceStatus::Blocked,
})
.await;
}
return Err(crate::error::OrchestratorError::PermissionStalled {
denied_path,
guidance,
});
}
Err(crate::error::OrchestratorError::PermissionBlocked {
denied_path,
guidance,
}) => {
use tracing::warn;
warn!(
"Permission auto-rejected for {} in workspace {}: {}",
change_id,
workspace_path.display(),
denied_path
);
if let Some(ref tx) = event_tx_for_permission {
use crate::parallel::ParallelEvent;
let _ = tx
.send(ParallelEvent::ApplyFailed {
change_id: change_id.to_string(),
error: format!("Permission auto-rejected: {}\n{}", denied_path, guidance),
})
.await;
}
return Err(crate::error::OrchestratorError::PermissionBlocked {
denied_path,
guidance,
});
}
Err(e) => return Err(e),
};
if let Some((attempt, findings)) = agent.get_acceptance_follow_up(change_id) {
acceptance_history
.lock()
.await
.set_follow_up_findings(change_id, attempt, findings);
}
mark_acceptance_context_injected(&agent, change_id, acceptance_tail_injected).await;
if apply_result.blocked_handoff.is_none() && apply_result.rejected_handoff.is_none() {
let (is_dirty, dirty_status) =
has_uncommitted_changes(workspace_path).await.map_err(|e| {
OrchestratorError::AgentCommand(format!(
"Failed to inspect worktree dirty state after apply completion for '{}': {}",
change_id, e
))
})?;
if is_dirty {
warn!(
change_id = %change_id,
workspace = %workspace_path.display(),
dirty_status = %dirty_status,
"Managed worktree is dirty after apply completion; running post-apply cleanup review before acceptance handoff"
);
run_post_apply_cleanup_review(
change_id,
workspace_path,
config,
ai_runner,
cancel_token,
event_tx_for_permission.as_ref(),
)
.await?;
}
}
info!(
"Apply completed for {} (revision={})",
change_id, apply_result.revision
);
Ok((
apply_result.revision,
apply_result.iterations,
apply_result.blocked_handoff,
apply_result.rejected_handoff,
))
}
#[allow(clippy::too_many_arguments)]
pub async fn execute_archive_finalization_in_workspace(
change_id: &str,
workspace_path: &Path,
config: &OrchestratorConfig,
event_tx: Option<mpsc::Sender<ParallelEvent>>,
vcs_backend: VcsBackend,
ai_runner: &AiCommandRunner,
shared_stagger_state: &Arc<Mutex<Option<std::time::Instant>>>,
) -> Result<String> {
use crate::execution::archive::{ensure_archive_commit, verify_archive_completion};
let verification = verify_archive_completion(change_id, Some(workspace_path));
if !verification.is_success() {
return Err(OrchestratorError::AgentCommand(format!(
"Cannot resume archive finalization for '{}': archive move has regressed; active change directory is present or archive files are missing in '{}'",
change_id,
workspace_path.display()
)));
}
if let Some(ref tx) = event_tx {
let _ = tx
.send(ParallelEvent::ArchiveResumed {
change_id: change_id.to_string(),
reason: Some("archive_commit_incomplete".to_string()),
summary: Some(
"Archive move is already complete; resuming commit finalization only"
.to_string(),
),
})
.await;
}
let resolve_agent =
AgentRunner::new_with_shared_state(config.clone(), shared_stagger_state.clone());
let change_id_owned = change_id.to_string();
let event_tx_clone = event_tx.clone();
ensure_archive_commit(
change_id,
workspace_path,
&resolve_agent,
ai_runner,
vcs_backend,
move |line| {
let event_tx = event_tx_clone.clone();
let change_id = change_id_owned.clone();
async move {
let text = match line {
OutputLine::Stdout(text) | OutputLine::Stderr(text) => text,
};
if let Some(ref tx) = event_tx {
if text.contains("Archive commit finalization retry scheduled") {
let _ = tx
.send(ParallelEvent::Log(
crate::events::LogEntry::warn(text.clone())
.with_change_id(&change_id)
.with_operation("archive-finalization"),
))
.await;
}
let _ = tx
.send(ParallelEvent::ArchiveOutput {
change_id,
output: text,
iteration: 1,
})
.await;
}
}
},
)
.await?;
if let Err(err) = delete_archive_state(workspace_path) {
warn!(
"Failed to delete archive state for {} after archive finalization resume: {}",
change_id, err
);
}
match vcs_backend {
VcsBackend::Git | VcsBackend::Auto => {
let revision = Command::new("git")
.args(["rev-parse", "HEAD"])
.current_dir(workspace_path)
.output()
.await
.map_err(|e| {
OrchestratorError::GitCommand(format!(
"Failed to get revision after archive finalization resume: {}",
e
))
})?;
if revision.status.success() {
Ok(String::from_utf8_lossy(&revision.stdout).trim().to_string())
} else {
Err(OrchestratorError::GitCommand(format!(
"Failed to get revision after archive finalization resume: {}",
String::from_utf8_lossy(&revision.stderr).trim()
)))
}
}
}
}
#[allow(clippy::too_many_arguments)]
pub async fn execute_archive_in_workspace(
change_id: &str,
workspace_path: &Path,
archive_cmd_template: &str,
config: &OrchestratorConfig,
event_tx: Option<mpsc::Sender<ParallelEvent>>,
vcs_backend: VcsBackend,
hooks: Option<&HookRunner>,
parallel_ctx: Option<&ParallelHookContext>,
cancel_token: Option<&CancellationToken>,
ai_runner: &AiCommandRunner,
archive_history: &Arc<Mutex<crate::history::ArchiveHistory>>,
apply_history: &Arc<Mutex<crate::history::ApplyHistory>>,
shared_stagger_state: &Arc<Mutex<Option<std::time::Instant>>>,
) -> Result<String> {
if cancel_token.is_some_and(|token| token.is_cancelled()) {
return Err(OrchestratorError::AgentCommand(format!(
"Cancelled archive for '{}' in workspace '{}'",
change_id,
workspace_path.display()
)));
}
use crate::execution::archive::get_task_progress;
let progress = match get_task_progress(change_id, Some(workspace_path)) {
Ok(Some(progress)) => {
if progress.total == 0 {
return Err(OrchestratorError::AgentCommand(format!(
"Cannot archive '{}' in workspace '{}': the task artifact exists but contains no tasks (0 tasks found)",
change_id,
workspace_path.display()
)));
}
if progress.completed < progress.total {
return Err(OrchestratorError::AgentCommand(format!(
"Cannot archive '{}' in workspace '{}': tasks not complete ({}/{} tasks completed)",
change_id,
workspace_path.display(),
progress.completed,
progress.total
)));
}
info!(
"Task verification passed for {}: {}/{} tasks completed",
change_id, progress.completed, progress.total
);
progress
}
Ok(None) => {
return Err(OrchestratorError::AgentCommand(format!(
"Cannot archive '{}' in workspace '{}': no task artifact ({} or {}) found under {} or in the archive directory",
change_id,
workspace_path.display(),
crate::task_file::MARKDOWN_FILE_NAME,
crate::task_file::JSON_FILE_NAME,
workspace_path
.join("openspec/changes")
.join(change_id)
.display()
)));
}
Err(e) => {
return Err(OrchestratorError::AgentCommand(format!(
"Cannot archive '{}' in workspace '{}': failed to read the task artifact: {}",
change_id,
workspace_path.display(),
e
)));
}
};
crate::vcs::git::commands::get_current_commit(workspace_path)
.await
.map_err(|e| {
OrchestratorError::AgentCommand(format!(
"Cannot archive '{}' in workspace '{}': failed to resolve current revision: {}",
change_id,
workspace_path.display(),
e
))
})?;
let stall_detector = StallDetector::new(config.get_stall_detection());
if let Some(hook_runner) = hooks {
let hook_ctx = build_parallel_hook_context(
change_id,
progress.completed,
progress.total,
0, parallel_ctx,
);
if let Some(ref tx) = event_tx {
let _ = tx
.send(ParallelEvent::HookStarted {
change_id: change_id.to_string(),
hook_type: "pre_archive".to_string(),
})
.await;
}
match hook_runner.run_hook(HookType::PreArchive, &hook_ctx).await {
Ok(()) => {
if let Some(ref tx) = event_tx {
let _ = tx
.send(ParallelEvent::HookCompleted {
change_id: change_id.to_string(),
hook_type: "pre_archive".to_string(),
})
.await;
}
}
Err(e) => {
error!("pre_archive hook failed for {}: {}", change_id, e);
if let Some(ref tx) = event_tx {
let _ = tx
.send(ParallelEvent::HookFailed {
change_id: change_id.to_string(),
hook_type: "pre_archive".to_string(),
error: e.to_string(),
})
.await;
}
return Err(e);
}
}
}
let user_prompt = config.get_archive_prompt();
let history_context = {
let history = archive_history.lock().await;
history.format_context(change_id)
};
let full_prompt = crate::agent::append_optional_prompt(
crate::agent::build_archive_prompt_with_skill(
config.get_archive_skill(),
Some(workspace_path),
change_id,
user_prompt,
&history_context,
),
config.get_archive_append_prompt(),
);
let command = OrchestratorConfig::expand_change_id(archive_cmd_template, change_id);
let command = OrchestratorConfig::expand_prompt(&command, &full_prompt);
debug!("Archive command in workspace: {}", command);
if let Some(ref tx) = event_tx {
let _ = tx
.send(ParallelEvent::ArchiveStarted {
change_id: change_id.to_string(),
command: command.clone(),
})
.await;
}
use crate::execution::archive::{
build_archive_error_message, ensure_archive_commit, extract_archive_runtime_blocker,
verify_archive_completion, ARCHIVE_COMMAND_MAX_RETRIES,
};
use crate::history::ArchivePrimaryReason;
let max_attempts = ARCHIVE_COMMAND_MAX_RETRIES.saturating_add(1);
let mut attempt: u32 = 0;
let is_git_repo = if matches!(vcs_backend, VcsBackend::Git | VcsBackend::Auto) {
match git_commands::check_git_repo(workspace_path).await {
Ok(is_repo) => is_repo,
Err(e) => {
warn!(
"Failed to check Git repository status for {}: {}",
change_id, e
);
false
}
}
} else {
false
};
let mut empty_commit_streak = 0u32;
loop {
attempt += 1;
let start = std::time::Instant::now();
debug!(
module = module_path!(),
"Executing shell command via AiCommandRunner with retry: {} (cwd: {:?})",
command,
workspace_path
);
let (mut child, mut output_rx) = ai_runner
.execute_streaming_with_retry(
&command,
Some(workspace_path),
Some("archive"),
Some(change_id),
)
.await?;
let mut output_collector = crate::history::OutputCollector::new();
use crate::ai_command_runner::OutputLine as AiOutputLine;
let change_id_clone = change_id.to_string();
let event_tx_clone = event_tx.clone();
loop {
let line = if let Some(token) = cancel_token {
tokio::select! {
_ = token.cancelled() => {
warn!(
change_id = %change_id,
workspace = %workspace_path.display(),
attempt = attempt,
"Archive cancellation observed while waiting for streaming output; terminating child"
);
let _ = child.terminate();
return Err(OrchestratorError::AgentCommand(format!(
"Cancelled archive for '{}' in workspace '{}'",
change_id,
workspace_path.display()
)));
}
line = output_rx.recv() => line,
}
} else {
output_rx.recv().await
};
let Some(line) = line else { break };
match &line {
AiOutputLine::Stdout(s) => output_collector.add_stdout(s),
AiOutputLine::Stderr(s) => output_collector.add_stderr(s),
}
if let Some(ref tx) = event_tx_clone {
let output_text = match line {
AiOutputLine::Stdout(s) | AiOutputLine::Stderr(s) => s,
};
let _ = tx
.send(ParallelEvent::ArchiveOutput {
change_id: change_id_clone.clone(),
output: output_text,
iteration: attempt,
})
.await;
}
}
let status = wait_for_streaming_child_with_cancel(
&mut child,
cancel_token,
"archive",
change_id,
workspace_path,
Some(attempt),
)
.await?;
if !status.success() {
return Err(OrchestratorError::AgentCommand(format!(
"Archive command failed for change '{}' in workspace '{}' (attempt {}) with exit code: {:?}",
change_id,
workspace_path.display(),
attempt,
status.code()
)));
}
if is_git_repo {
if let Err(e) =
git_commands::create_archive_wip_commit(workspace_path, change_id, attempt).await
{
warn!(
"Failed to create WIP(archive) commit for {} (attempt {}): {}",
change_id, attempt, e
);
} else if stall_detector.config().enabled {
match git_commands::is_head_empty_commit(workspace_path).await {
Ok(is_empty) => {
if is_empty {
empty_commit_streak = empty_commit_streak.saturating_add(1);
} else {
empty_commit_streak = 0;
}
if empty_commit_streak >= stall_detector.config().threshold {
let message = format!(
"Stall detected for {} after {} empty WIP commits (archive)",
change_id, empty_commit_streak
);
warn!(
"{} (threshold {})",
message,
stall_detector.config().threshold
);
return Err(OrchestratorError::AgentCommand(message));
}
}
Err(e) => {
warn!(
"Failed to check WIP(archive) commit for {} (attempt {}): {}",
change_id, attempt, e
);
}
}
}
}
let verification = verify_archive_completion(change_id, Some(workspace_path));
{
let mut history = archive_history.lock().await;
let verification_result = if verification.is_success() {
None
} else {
Some(format!(
"Change still exists at openspec/changes/{}",
change_id
))
};
let attempt_record = crate::history::ArchiveAttempt {
attempt: history.count(change_id) + 1,
success: status.success() && verification.is_success(),
duration: start.elapsed(),
error: if status.success() && verification.is_success() {
None
} else if !status.success() {
Some(format!("Exit code: {:?}", status.code()))
} else {
Some("Archive command succeeded but verification failed".to_string())
},
primary_reason: if status.success() && verification.is_success() {
None
} else if !status.success() {
Some(ArchivePrimaryReason::CommandFailed)
} else {
Some(ArchivePrimaryReason::VerificationFailed)
},
verification_result,
exit_code: status.code(),
stdout_tail: output_collector.stdout_tail(),
stderr_tail: output_collector.stderr_tail(),
};
history.record(change_id, attempt_record);
}
if verification.is_success() {
break;
}
if attempt <= ARCHIVE_COMMAND_MAX_RETRIES {
let retry_summary =
"archive verification failed; change directory still exists".to_string();
if let Some(ref tx) = event_tx {
let _ = tx
.send(ParallelEvent::ArchiveRetryScheduled {
change_id: change_id.to_string(),
attempt,
max_attempts,
reason: Some(
ArchivePrimaryReason::VerificationFailed
.as_str()
.to_string(),
),
summary: Some(retry_summary),
})
.await;
let _ = tx
.send(ParallelEvent::Log(
crate::events::LogEntry::warn(format!(
"Archive verification failed for {} (attempt {}/{}); retrying archive command",
change_id, attempt, max_attempts
))
.with_change_id(change_id)
.with_operation("archive")
.with_iteration(attempt),
))
.await;
}
warn!(
change_id = %change_id,
attempt = attempt,
max_attempts = max_attempts,
"Archive verification failed; retrying archive command"
);
continue;
}
let runtime_blocker = extract_archive_runtime_blocker(
output_collector.stdout_tail().as_deref(),
output_collector.stderr_tail().as_deref(),
);
let final_error = build_archive_error_message(
change_id,
Some(workspace_path),
runtime_blocker.as_deref(),
);
return Err(OrchestratorError::AgentCommand(final_error));
}
info!(
"Archive verification passed for {}: change moved to archive",
change_id
);
if is_git_repo {
if let Err(e) = git_commands::squash_archive_wip_commits(workspace_path, change_id).await {
warn!(
"Failed to squash WIP(archive) commits for {}: {}",
change_id, e
);
}
}
let resolve_agent =
AgentRunner::new_with_shared_state(config.clone(), shared_stagger_state.clone());
let change_id_owned = change_id.to_string();
let event_tx_clone = event_tx.clone();
let final_attempt = attempt;
ensure_archive_commit(
change_id,
workspace_path,
&resolve_agent,
ai_runner,
vcs_backend,
move |line| {
let event_tx = event_tx_clone.clone();
let change_id = change_id_owned.clone();
let iteration = final_attempt;
async move {
let text = match line {
OutputLine::Stdout(text) | OutputLine::Stderr(text) => text,
};
if let Some(ref tx) = event_tx {
if text.contains("Archive commit finalization retry scheduled") {
let _ = tx
.send(ParallelEvent::Log(
crate::events::LogEntry::warn(text.clone())
.with_change_id(&change_id)
.with_operation("archive-finalization")
.with_iteration(iteration),
))
.await;
}
let _ = tx
.send(ParallelEvent::ArchiveOutput {
change_id,
output: text,
iteration,
})
.await;
}
}
},
)
.await?;
let revision = match vcs_backend {
VcsBackend::Git | VcsBackend::Auto => {
debug!(
module = module_path!(),
"Executing git command: git rev-parse HEAD (cwd: {:?})", workspace_path
);
if !workspace_path.exists() {
warn!(
"Workspace path {:?} no longer exists after archive (likely deleted by archive command), using placeholder revision",
workspace_path
);
"archived".to_string()
} else {
match Command::new("git")
.args(["rev-parse", "HEAD"])
.current_dir(workspace_path)
.output()
.await
{
Ok(revision_output) if revision_output.status.success() => {
String::from_utf8_lossy(&revision_output.stdout)
.trim()
.to_string()
}
Ok(revision_output) => {
let stderr = String::from_utf8_lossy(&revision_output.stderr);
warn!(
"Failed to get revision from workspace {:?} after archive: {} (likely deleted by archive command), using placeholder",
workspace_path, stderr
);
"archived".to_string()
}
Err(e) => {
warn!(
"Failed to execute git rev-parse in workspace {:?} after archive: {} (likely deleted by archive command), using placeholder",
workspace_path, e
);
"archived".to_string()
}
}
}
}
};
if let Some(hook_runner) = hooks {
let hook_ctx = build_parallel_hook_context(
change_id,
progress.completed,
progress.total,
0, parallel_ctx,
);
if let Some(ref tx) = event_tx {
let _ = tx
.send(ParallelEvent::HookStarted {
change_id: change_id.to_string(),
hook_type: "post_archive".to_string(),
})
.await;
}
match hook_runner.run_hook(HookType::PostArchive, &hook_ctx).await {
Ok(()) => {
if let Some(ref tx) = event_tx {
let _ = tx
.send(ParallelEvent::HookCompleted {
change_id: change_id.to_string(),
hook_type: "post_archive".to_string(),
})
.await;
}
}
Err(e) => {
error!("post_archive hook failed for {}: {}", change_id, e);
if let Some(ref tx) = event_tx {
let _ = tx
.send(ParallelEvent::HookFailed {
change_id: change_id.to_string(),
hook_type: "post_archive".to_string(),
error: e.to_string(),
})
.await;
}
return Err(e);
}
}
}
{
let mut apply_hist = apply_history.lock().await;
apply_hist.clear(change_id);
let mut archive_hist = archive_history.lock().await;
archive_hist.clear(change_id);
}
clear_acceptance_state(change_id, workspace_path, config).await;
if let Err(err) = delete_archive_state(workspace_path) {
warn!(
"Failed to delete archive state for {} after archive completion: {}",
change_id, err
);
}
Ok(revision)
}
fn format_acceptance_failure_log_message(
findings: &[crate::acceptance::AcceptanceFinding],
) -> String {
let finding_count = findings.len();
let blocking_gate_context = findings
.first()
.map(|finding| finding.text().to_string())
.unwrap_or_else(|| "no acceptance findings captured".to_string());
format!(
"Acceptance failed ({} findings), blocking gate context: {}",
finding_count, blocking_gate_context
)
}
fn resolve_acceptance_state_revision(start_revision: &str, end_revision: Option<String>) -> String {
end_revision.unwrap_or_else(|| start_revision.to_string())
}
fn revision_to_history_commit_hash(revision: &str) -> Option<String> {
if revision == "unknown" {
None
} else {
Some(revision.to_string())
}
}
const ACCEPTANCE_VERDICT_GRACE_DEFAULT_SECS: u64 = 30;
tokio::task_local! {
pub(crate) static VERDICT_GRACE_OVERRIDE_SECS: u64;
}
pub(crate) fn acceptance_verdict_grace_period() -> std::time::Duration {
let secs = VERDICT_GRACE_OVERRIDE_SECS
.try_with(|secs| *secs)
.ok()
.filter(|secs| *secs > 0)
.unwrap_or(ACCEPTANCE_VERDICT_GRACE_DEFAULT_SECS);
std::time::Duration::from_secs(secs)
}
#[cfg(test)]
pub(crate) async fn scoped_verdict_grace_secs_for_test<F, R>(secs: u64, fut: F) -> R
where
F: std::future::Future<Output = R>,
{
VERDICT_GRACE_OVERRIDE_SECS.scope(secs, fut).await
}
tokio::task_local! {
pub(crate) static REVIEW_BUDGET_OVERRIDE_MS: u64;
}
#[cfg(test)]
pub(crate) async fn scoped_review_budget_millis_for_test<F, R>(millis: u64, fut: F) -> R
where
F: std::future::Future<Output = R>,
{
REVIEW_BUDGET_OVERRIDE_MS.scope(millis, fut).await
}
fn effective_review_budget(derived: Option<std::time::Duration>) -> Option<std::time::Duration> {
REVIEW_BUDGET_OVERRIDE_MS
.try_with(|millis| *millis)
.ok()
.filter(|millis| *millis > 0)
.map(std::time::Duration::from_millis)
.or(derived)
}
pub(crate) struct AcceptanceBoundary {
pub manifest: crate::orchestration::acceptance::execution_manifest::AcceptanceExecutionManifest,
pub gates: crate::orchestration::acceptance::execution_manifest::GateResults,
pub verdict: crate::orchestration::acceptance::gate_execution::GateVerdict,
pub review_budget: Option<std::time::Duration>,
pub store: crate::orchestration::acceptance::evidence_location::AcceptanceStore,
pub review_base_ref: Option<String>,
}
pub(crate) async fn run_acceptance_boundary(
change_id: &str,
workspace_path: &Path,
skill_name: &str,
review_base: Option<&str>,
absolute_deadline_secs: u64,
state_base_dir: Option<&str>,
) -> std::result::Result<
Option<AcceptanceBoundary>,
crate::orchestration::acceptance::execution_manifest::AcceptanceExecutionHold,
> {
use crate::orchestration::acceptance::evidence_location::prepare_store;
use crate::orchestration::acceptance::execution_manifest::{
live_acceptance_holds, partition_deadline, AcceptanceExecutionHold, AcceptanceHoldCategory,
};
use crate::orchestration::acceptance::gate_execution::{
classify_gate_results, execute_declared_gates, plan_manifest_reuse, GateBudget,
};
use crate::orchestration::acceptance::manifest_builder::{
build_manifest_for_workspace, repository_project_root,
};
use crate::orchestration::acceptance::verification_evidence::{
DirectCommandSupervisor, GitRepositoryFacts, ReusePolicy, RuntimeVerificationExecutor,
SystemClock,
};
let partition = partition_deadline(absolute_deadline_secs);
if !crate::orchestration::acceptance::manifest_builder::workspace_is_repository(workspace_path)
.await
{
debug!(
change_id = change_id,
workspace = %workspace_path.display(),
"Acceptance boundary skipped: the workspace is not a Git working tree"
);
return Ok(None);
}
let project_root = repository_project_root(workspace_path)
.await
.unwrap_or_else(|| workspace_path.to_path_buf());
let store = match prepare_store(state_base_dir, &project_root, workspace_path, change_id) {
Ok(store) => store,
Err(refusal) => {
let category = refusal.category();
warn!(
change_id = change_id,
category = category.as_str(),
"Acceptance not dispatched: {}",
refusal.evidence()
);
return Err(AcceptanceExecutionHold {
category,
budget_secs: partition.absolute_secs,
cleanup_confirmed: true,
cleanup_diagnostics: "no gate and no reviewer was started".to_string(),
fingerprint: "0".repeat(64),
evidence: vec![refusal.evidence()],
});
}
};
let evidence_store = store.evidence();
let manifest = match build_manifest_for_workspace(
&GitRepositoryFacts,
workspace_path,
&evidence_store,
change_id,
skill_name,
review_base,
absolute_deadline_secs,
)
.await
{
Ok(manifest) => manifest,
Err(error) => {
return Err(AcceptanceExecutionHold {
category: AcceptanceHoldCategory::ManifestMalformed,
budget_secs: partition.absolute_secs,
cleanup_confirmed: true,
cleanup_diagnostics: "no reviewer was started".to_string(),
fingerprint: "0".repeat(64),
evidence: vec![error.detail()],
});
}
};
let fingerprint = manifest.fingerprint();
let holds = live_acceptance_holds();
let admission = holds.classify(change_id, &fingerprint);
if let crate::orchestration::acceptance::execution_manifest::AcceptanceAdmission::Refuse {
category,
..
} = &admission
{
warn!(
change_id = change_id,
category = category.as_str(),
"Acceptance not dispatched: input is unchanged since this owner's non-resumable hold"
);
return Err(AcceptanceExecutionHold {
category: *category,
budget_secs: partition.absolute_secs,
cleanup_confirmed: true,
cleanup_diagnostics: "no reviewer was started".to_string(),
fingerprint,
evidence: admission.detail().into_iter().collect(),
});
}
holds.clear(change_id);
let manifest_store = store.manifests();
if let Err(error) = manifest_store.store(&manifest) {
warn!(
change_id = change_id,
"Acceptance manifest could not be persisted to the external store: {}", error
);
}
let reuse_decisions = plan_manifest_reuse(
&GitRepositoryFacts,
workspace_path,
&evidence_store,
&manifest,
ReusePolicy::default(),
)
.await;
let gate_budget = match partition.work_duration() {
Some(work) => GateBudget::bounded(work),
None => GateBudget::unbounded(),
};
let gate_start = std::time::Instant::now();
let executor = RuntimeVerificationExecutor::new(
workspace_path.to_path_buf(),
evidence_store.clone(),
GitRepositoryFacts,
DirectCommandSupervisor,
SystemClock,
)
.with_review_binding(
crate::orchestration::acceptance::verification_evidence::ReviewBinding {
base_commit: manifest.review_base_commit.clone(),
range: manifest.review_range.clone(),
},
);
let gates = execute_declared_gates(&executor, &manifest, &reuse_decisions, gate_budget).await;
for outcome in &gates.outcomes {
info!(
change_id = change_id,
verification_id = outcome.verification_id(),
status = outcome.status(),
"Declared Acceptance gate: {}",
outcome.summary()
);
}
let review_budget = partition
.work_duration()
.map(|work| work.saturating_sub(gate_start.elapsed()));
let mut manifest = manifest;
{
use crate::orchestration::acceptance::verification_evidence::RepositoryFacts;
for gate in &mut manifest.gates {
gate.artifact_digest = GitRepositoryFacts
.hash_file(
workspace_path,
&evidence_store.artifact_path(&gate.verification_id),
)
.await
.ok();
}
}
let settled_fingerprint = manifest.fingerprint();
if settled_fingerprint != fingerprint {
if let Err(error) = manifest_store.store(&manifest) {
warn!(
change_id = change_id,
"Settled Acceptance manifest could not be persisted: {}", error
);
}
}
debug!(
change_id = change_id,
fingerprint = %settled_fingerprint,
pre_gate_fingerprint = %fingerprint,
store = %store.root().display(),
"Acceptance execution manifest frozen"
);
Ok(Some(AcceptanceBoundary {
verdict: classify_gate_results(&gates),
manifest,
gates,
review_budget,
store,
review_base_ref: review_base.map(str::to_string),
}))
}
pub(crate) async fn clear_acceptance_state(
change_id: &str,
workspace_path: &Path,
config: &OrchestratorConfig,
) {
use crate::orchestration::acceptance::evidence_location::{store_path, AcceptanceStore};
use crate::orchestration::acceptance::execution_manifest::live_acceptance_holds;
use crate::orchestration::acceptance::manifest_builder::repository_project_root;
live_acceptance_holds().clear(change_id);
let project_root = repository_project_root(workspace_path)
.await
.unwrap_or_else(|| workspace_path.to_path_buf());
match store_path(
config.get_state_base_dir(),
&project_root,
workspace_path,
change_id,
) {
Ok(root) => AcceptanceStore::at(root).clear(),
Err(error) => debug!(
change_id = change_id,
"Acceptance cache was not removed because its path did not resolve: {}", error
),
}
}
pub(crate) fn record_acceptance_hold(
boundary: &AcceptanceBoundary,
hold: &crate::orchestration::acceptance::execution_manifest::AcceptanceExecutionHold,
) {
use crate::orchestration::acceptance::execution_manifest::{
live_acceptance_holds, AcceptanceDiagnostics, LiveAcceptanceHold,
};
let change_id = boundary.manifest.change_id.as_str();
live_acceptance_holds().record(
change_id,
LiveAcceptanceHold {
category: hold.category,
fingerprint: hold.fingerprint.clone(),
review_base_ref: boundary.review_base_ref.clone(),
},
);
if let Err(error) = boundary
.store
.manifests()
.record_diagnostics(&AcceptanceDiagnostics {
hold_category: Some(hold.category),
fingerprint: hold.fingerprint.clone(),
cleanup_confirmed: hold.cleanup_confirmed,
cleanup_diagnostics: hold.cleanup_diagnostics.clone(),
evidence: hold.evidence.clone(),
recorded_at: chrono::Utc::now(),
recorded_by_pid: std::process::id(),
})
{
warn!(
change_id = %change_id,
"Acceptance diagnostics could not be persisted to the external store: {}", error
);
}
}
pub(crate) async fn emit_acceptance_execution_hold(
change_id: &str,
hold: &crate::orchestration::acceptance::execution_manifest::AcceptanceExecutionHold,
event_tx: Option<&mpsc::Sender<ParallelEvent>>,
iteration: u32,
) {
if let Some(tx) = event_tx {
let _ = tx
.send(ParallelEvent::Log(
crate::events::LogEntry::error(hold.summary(change_id))
.with_change_id(change_id)
.with_operation("acceptance")
.with_iteration(iteration),
))
.await;
let _ = tx
.send(ParallelEvent::AcceptanceCompleted {
change_id: change_id.to_string(),
})
.await;
}
}
#[allow(clippy::too_many_arguments)]
pub async fn execute_acceptance_in_workspace(
change_id: &str,
workspace_path: &Path,
agent: &mut AgentRunner,
event_tx: Option<mpsc::Sender<ParallelEvent>>,
cancel_token: Option<&CancellationToken>,
ai_runner: &AiCommandRunner,
config: &OrchestratorConfig,
acceptance_tail_injected: &Arc<Mutex<std::collections::HashMap<String, bool>>>,
acceptance_history: &Arc<Mutex<crate::history::AcceptanceHistory>>,
base_branch: Option<&str>,
protocol_retry: Option<crate::orchestration::acceptance::AcceptanceProtocolRetry>,
command_mode: crate::orchestration::acceptance::AcceptanceCommandMode,
) -> Result<(crate::orchestration::AcceptanceResult, u32)> {
use crate::acceptance::{parse_acceptance_output, AcceptanceResult as ParseResult};
if cancel_token.is_some_and(|token| token.is_cancelled()) {
return Ok((crate::orchestration::AcceptanceResult::Cancelled, 0));
}
info!("Running acceptance test for {} in workspace", change_id);
let start_time_for_gate_phase = std::time::Instant::now();
let acceptance_runtime_limit_secs = ai_runner
.queue_config()
.effective_max_runtime_secs(Some(crate::command_queue::ACCEPTANCE_OPERATION_TYPE));
let commit_hash = crate::vcs::git::commands::get_current_commit(workspace_path)
.await
.ok();
let acceptance_iteration = agent.next_acceptance_attempt_number(change_id);
let user_prompt = config.get_acceptance_prompt();
let history_context = agent.format_acceptance_history(change_id);
let diff_context = {
let current_commit = crate::vcs::git::commands::get_current_commit(workspace_path)
.await
.ok();
let base_commit = {
let acc_history = acceptance_history.lock().await;
if acc_history.count(change_id) == 0 {
base_branch.map(|b| b.to_string())
} else {
acc_history.last_commit_hash(change_id)
}
};
if let (Some(base), Some(current)) = (base_commit.as_ref(), current_commit.as_ref()) {
match crate::vcs::git::commands::get_changed_files(workspace_path, Some(base), current)
.await
{
Ok(files) => {
let previous_findings = {
let acc_history = acceptance_history.lock().await;
if acc_history.count(change_id) > 0 {
acc_history
.last_findings(change_id)
.map(|findings| crate::acceptance::finding_texts(&findings))
} else {
None
}
};
if !files.is_empty() || previous_findings.is_some() {
crate::agent::build_acceptance_diff_context(
&files,
previous_findings.as_deref(),
)
} else {
String::new()
}
}
Err(e) => {
warn!(
"Failed to get changed files for acceptance diff context: {}",
e
);
String::new()
}
}
} else {
String::new()
}
};
let stdout_tail = agent.get_last_acceptance_stdout_tail(change_id);
let stderr_tail = agent.get_last_acceptance_stderr_tail(change_id);
let previous_findings_text = agent
.get_last_acceptance_finding_texts(change_id)
.map(|findings| findings.join("\n"));
let previous_acceptance_denial = agent
.acceptance_command_recovery(change_id)
.and_then(crate::orchestration::acceptance::classify_acceptance_command_denial)
.or_else(|| {
crate::permission::classify_permission_denial(&[
stdout_tail.as_deref(),
stderr_tail.as_deref(),
previous_findings_text.as_deref(),
])
});
let last_output_context = crate::agent::build_last_acceptance_output_context(
stdout_tail.as_deref(),
stderr_tail.as_deref(),
);
let protocol_retry_context = protocol_retry.map_or_else(String::new, |retry| {
crate::agent::build_missing_verdict_continuation_context(
retry,
stdout_tail.as_deref(),
stderr_tail.as_deref(),
agent
.get_last_acceptance_finding_texts(change_id)
.as_deref(),
)
});
let command_recovery_context = crate::agent::build_acceptance_command_recovery_context(
agent.acceptance_command_recovery(change_id),
);
let boundary = match run_acceptance_boundary(
change_id,
workspace_path,
config.get_accept_skill(),
base_branch,
acceptance_runtime_limit_secs,
config.get_state_base_dir(),
)
.await
{
Ok(boundary) => boundary,
Err(hold) => {
warn!(
change_id = %change_id,
category = hold.category.as_str(),
"Acceptance boundary could not be derived: {}",
hold.summary(change_id)
);
emit_acceptance_execution_hold(change_id, &hold, event_tx.as_ref(), 0).await;
return Ok((
crate::orchestration::AcceptanceResult::ExecutionHold { hold },
0,
));
}
};
if let Some(boundary) = boundary.as_ref() {
use crate::orchestration::acceptance::gate_execution::GateVerdict;
match &boundary.verdict {
GateVerdict::Proceed => {}
GateVerdict::Failed { evidence } => {
let attempt_number = agent.next_acceptance_attempt_number(change_id);
let mut findings = vec![format!(
"Declared change-blocking verification failed under runtime supervision \
against candidate {}.",
&boundary.manifest.candidate_commit_oid
[..12.min(boundary.manifest.candidate_commit_oid.len())]
)];
findings.extend(evidence.iter().cloned());
let findings = crate::acceptance::legacy_findings(findings);
let attempt = crate::history::AcceptanceAttempt {
attempt: attempt_number,
passed: false,
duration: start_time_for_gate_phase.elapsed(),
findings: Some(findings.clone()),
exit_code: None,
stdout_tail: None,
stderr_tail: None,
commit_hash: Some(boundary.manifest.candidate_commit_oid.clone()),
};
agent.record_acceptance_attempt(change_id, attempt.clone());
acceptance_history.lock().await.record(change_id, attempt);
acceptance_tail_injected.lock().await.remove(change_id);
if let Some(ref tx) = event_tx {
let _ = tx
.send(ParallelEvent::Log(
crate::events::LogEntry::error(format!(
"Acceptance failed on a declared change-blocking verification: {}",
evidence.join(" | ")
))
.with_change_id(change_id)
.with_operation("acceptance")
.with_iteration(attempt_number),
))
.await;
let _ = tx
.send(ParallelEvent::AcceptanceCompleted {
change_id: change_id.to_string(),
})
.await;
}
return Ok((
crate::orchestration::AcceptanceResult::Fail { findings },
attempt_number,
));
}
GateVerdict::ExternalPrerequisite { blocker } => {
let attempt_number = agent.next_acceptance_attempt_number(change_id);
if let Some(ref tx) = event_tx {
let _ = tx
.send(ParallelEvent::Log(
crate::events::LogEntry::warn(format!(
"Acceptance blocked on an external prerequisite for a declared \
verification ({}): {}",
blocker.category, blocker.next_action
))
.with_change_id(change_id)
.with_operation("acceptance")
.with_iteration(attempt_number),
))
.await;
let _ = tx
.send(ParallelEvent::AcceptanceCompleted {
change_id: change_id.to_string(),
})
.await;
}
return Ok((
crate::orchestration::AcceptanceResult::Stalled {
blocker: blocker.clone(),
},
attempt_number,
));
}
GateVerdict::Hold { category, evidence } => {
let (cleanup_confirmed, cleanup_diagnostics) = boundary.gates.cleanup_evidence();
let hold =
crate::orchestration::acceptance::execution_manifest::AcceptanceExecutionHold {
category: *category,
budget_secs: acceptance_runtime_limit_secs,
cleanup_confirmed,
cleanup_diagnostics,
fingerprint: boundary.manifest.fingerprint(),
evidence: evidence.clone(),
};
record_acceptance_hold(boundary, &hold);
emit_acceptance_execution_hold(change_id, &hold, event_tx.as_ref(), 0).await;
return Ok((
crate::orchestration::AcceptanceResult::ExecutionHold { hold },
0,
));
}
}
}
let boundary_context = boundary.as_ref().map_or_else(String::new, |boundary| {
crate::agent::build_acceptance_boundary_context(
&boundary.manifest,
&boundary.gates,
boundary.store.root(),
)
});
let full_prompt = crate::agent::append_optional_prompt(
match config.get_acceptance_prompt_mode() {
crate::config::AcceptancePromptMode::Full => {
crate::agent::build_acceptance_prompt_with_skill(
config.get_accept_skill(),
Some(workspace_path),
change_id,
user_prompt,
&history_context,
&last_output_context,
&diff_context,
&protocol_retry_context,
&command_recovery_context,
)
}
crate::config::AcceptancePromptMode::ContextOnly => {
crate::agent::build_acceptance_prompt_context_only_with_skill(
config.get_accept_skill(),
Some(workspace_path),
change_id,
user_prompt,
&history_context,
&last_output_context,
&diff_context,
&protocol_retry_context,
&command_recovery_context,
)
}
},
config.get_acceptance_append_prompt(),
);
let full_prompt = crate::agent::append_optional_prompt(full_prompt, Some(&boundary_context));
let full_prompt = crate::agent::append_optional_prompt(
full_prompt,
Some(&crate::agent::build_acceptance_escalation_context(
command_mode,
)),
);
let template =
crate::orchestration::acceptance::acceptance_command_template(config, command_mode)?;
let command = OrchestratorConfig::expand_change_id(template, change_id);
let command = OrchestratorConfig::expand_prompt(&command, &full_prompt);
debug!(
module = module_path!(),
command = %crate::events::command_log_summary(&command),
command_mode = command_mode.label(),
cwd = ?workspace_path,
"Executing acceptance command via AiCommandRunner"
);
let start_revision = commit_hash.clone().unwrap_or_else(|| "unknown".to_string());
if let Some(ref tx) = event_tx {
let _ = tx
.send(ParallelEvent::AcceptanceStarted {
change_id: change_id.to_string(),
command: command.clone(),
})
.await;
}
let start_time = std::time::Instant::now();
let (mut child, mut output_rx) = ai_runner
.execute_streaming_with_retry(
&command,
Some(workspace_path),
Some(crate::command_queue::ACCEPTANCE_OPERATION_TYPE),
Some(change_id),
)
.await?;
let mut output_collector = crate::history::OutputCollector::new();
let mut full_stdout = String::new();
let verdict_grace_period = acceptance_verdict_grace_period();
let mut verdict_detected = false;
let mut verdict_stream_detector = crate::acceptance::VerdictStreamDetector::default();
let mut verdict_deadline: Option<tokio::time::Instant> = None;
let mut early_terminated = false;
let review_deadline: Option<tokio::time::Instant> = boundary
.as_ref()
.and_then(|boundary| effective_review_budget(boundary.review_budget))
.map(|budget| tokio::time::Instant::now() + budget);
let mut review_deadline_expired = false;
use crate::ai_command_runner::OutputLine as AiOutputLine;
loop {
let active_deadline = match (verdict_deadline, review_deadline) {
(Some(verdict), Some(review)) => Some((verdict.min(review), review < verdict)),
(Some(verdict), None) => Some((verdict, false)),
(None, Some(review)) => Some((review, true)),
(None, None) => None,
};
let line = if let Some((deadline, is_review_deadline)) = active_deadline {
let recv_future = output_rx.recv();
let recv_with_deadline = tokio::time::timeout_at(deadline, recv_future);
let recv_result = if let Some(token) = cancel_token {
tokio::select! {
_ = token.cancelled() => {
warn!(
change_id = %change_id,
workspace = %workspace_path.display(),
iteration = acceptance_iteration,
"Acceptance cancellation observed while waiting for streaming output; terminating child"
);
let _ = child.terminate();
return Ok((crate::orchestration::AcceptanceResult::Cancelled, 0));
}
result = recv_with_deadline => result,
}
} else {
recv_with_deadline.await
};
match recv_result {
Ok(Some(line)) => line,
Ok(None) => break,
Err(_) if is_review_deadline => {
warn!(
change_id = %change_id,
workspace = %workspace_path.display(),
iteration = acceptance_iteration,
"Acceptance review deadline reached; terminating the owned reviewer \
process group inside the reserved finalization window"
);
let _ = child.terminate();
early_terminated = true;
review_deadline_expired = true;
break;
}
Err(_) => {
info!(
"Acceptance verdict grace period ({}s) expired for {}, terminating child process",
verdict_grace_period.as_secs(),
change_id
);
let _ = child.terminate();
early_terminated = true;
break;
}
}
} else if let Some(token) = cancel_token {
tokio::select! {
_ = token.cancelled() => {
warn!(
change_id = %change_id,
workspace = %workspace_path.display(),
iteration = acceptance_iteration,
"Acceptance cancellation observed while waiting for streaming output; terminating child"
);
let _ = child.terminate();
return Ok((crate::orchestration::AcceptanceResult::Cancelled, 0));
}
line = output_rx.recv() => match line {
Some(line) => line,
None => break,
},
}
} else {
match output_rx.recv().await {
Some(line) => line,
None => break,
}
};
if cancel_token.is_some_and(|token| token.is_cancelled()) {
warn!("Acceptance test cancelled for: {}", change_id);
let _ = child.terminate();
return Ok((crate::orchestration::AcceptanceResult::Cancelled, 0));
}
match line {
AiOutputLine::Stdout(s) => {
output_collector.add_stdout(&s);
full_stdout.push_str(&s);
full_stdout.push('\n');
if let Some(ref tx) = event_tx {
let _ = tx
.send(ParallelEvent::Log(
crate::events::LogEntry::info(&s)
.with_change_id(change_id)
.with_operation("acceptance")
.with_iteration(acceptance_iteration),
))
.await;
}
if !verdict_detected && verdict_stream_detector.detect(&s).is_some() {
verdict_detected = true;
verdict_deadline = Some(tokio::time::Instant::now() + verdict_grace_period);
info!(
"Acceptance canonical verdict detected for {}, starting {}s grace period",
change_id,
verdict_grace_period.as_secs()
);
}
}
AiOutputLine::Stderr(s) => {
output_collector.add_stderr(&s);
if let Some(ref tx) = event_tx {
let _ = tx
.send(ParallelEvent::Log(
crate::events::LogEntry::info(&s)
.with_change_id(change_id)
.with_operation("acceptance")
.with_iteration(acceptance_iteration),
))
.await;
}
}
}
}
let status = match wait_for_streaming_child_with_cancel(
&mut child,
cancel_token,
"acceptance",
change_id,
workspace_path,
Some(acceptance_iteration),
)
.await
{
Ok(status) => status,
Err(err) if cancel_token.is_some_and(|token| token.is_cancelled()) => {
warn!(
change_id = %change_id,
workspace = %workspace_path.display(),
"Acceptance cancelled while waiting for child status: {}",
err
);
return Ok((crate::orchestration::AcceptanceResult::Cancelled, 0));
}
Err(err) => return Err(err),
};
if let (true, false, Some(boundary)) =
(review_deadline_expired, verdict_detected, boundary.as_ref())
{
let cleanup = child.process_group_cleanup().await;
let category = if cleanup.is_confirmed() {
crate::orchestration::acceptance::execution_manifest::AcceptanceHoldCategory::ReviewDeadlineExhausted
} else {
crate::orchestration::acceptance::execution_manifest::AcceptanceHoldCategory::ReviewCleanupUnproven
};
let hold = crate::orchestration::acceptance::execution_manifest::AcceptanceExecutionHold {
category,
budget_secs: acceptance_runtime_limit_secs,
cleanup_confirmed: cleanup.is_confirmed(),
cleanup_diagnostics: cleanup.diagnostics(),
fingerprint: boundary.manifest.fingerprint(),
evidence: {
let mut evidence = boundary.gates.summary_lines();
evidence.push(
"semantic review emitted no canonical verdict before the reserved \
finalization window"
.to_string(),
);
evidence
},
};
warn!(
change_id = %change_id,
workspace = %workspace_path.display(),
iteration = acceptance_iteration,
category = hold.category.as_str(),
cleanup_confirmed = hold.cleanup_confirmed,
"Acceptance review deadline exhausted: {}",
hold.cleanup_diagnostics
);
record_acceptance_hold(boundary, &hold);
emit_acceptance_execution_hold(change_id, &hold, event_tx.as_ref(), acceptance_iteration)
.await;
return Ok((
crate::orchestration::AcceptanceResult::ExecutionHold { hold },
0,
));
}
{
let termination = child.termination().await;
if termination.is_runtime_limit() {
let cleanup = child.process_group_cleanup().await;
let limit = crate::orchestration::acceptance::classify_acceptance_runtime_limit(
termination,
acceptance_runtime_limit_secs,
&cleanup,
)
.expect("runtime-limit termination always classifies as a runtime limit");
warn!(
change_id = %change_id,
workspace = %workspace_path.display(),
iteration = acceptance_iteration,
limit_secs = limit.limit_secs,
cleanup_confirmed = limit.cleanup_confirmed,
"Acceptance exceeded its absolute runtime limit: {}",
limit.cleanup_diagnostics
);
if let Some(ref tx) = event_tx {
let _ = tx
.send(ParallelEvent::Log(
crate::events::LogEntry::error(limit.summary(change_id))
.with_change_id(change_id)
.with_operation("acceptance")
.with_iteration(acceptance_iteration),
))
.await;
let _ = tx
.send(ParallelEvent::AcceptanceCompleted {
change_id: change_id.to_string(),
})
.await;
}
return Ok((
crate::orchestration::AcceptanceResult::RuntimeLimit { limit },
0,
));
}
}
let verdict_finalized_run = early_terminated && verdict_detected;
let end_revision = resolve_acceptance_state_revision(
&start_revision,
crate::vcs::git::commands::get_current_commit(workspace_path)
.await
.ok(),
);
if end_revision != start_revision {
warn!(
module = module_path!(),
change_id = %change_id,
start_revision = %start_revision,
end_revision = %end_revision,
workspace = %workspace_path.display(),
"Acceptance updated HEAD during execution; durable acceptance state will use end revision"
);
}
let stdout_tail = output_collector.stdout_tail();
let stderr_tail = output_collector.stderr_tail();
let parse_result = parse_acceptance_output(&full_stdout);
let tail_findings = build_acceptance_tail_findings(stdout_tail.clone(), stderr_tail.clone());
if !status.success() && !verdict_finalized_run {
let current_denial = crate::permission::classify_permission_denial(&[
stdout_tail.as_deref(),
stderr_tail.as_deref(),
]);
if let Some(denial) = ¤t_denial {
let repeated_unresolved = previous_acceptance_denial
.as_ref()
.is_some_and(|previous| previous.signature() == denial.signature())
&& end_revision == start_revision;
if repeated_unresolved {
let attempt_number = agent.next_acceptance_attempt_number(change_id);
let attempt = crate::history::AcceptanceAttempt {
attempt: attempt_number,
passed: false,
duration: start_time.elapsed(),
findings: Some(crate::acceptance::legacy_findings(tail_findings.clone())),
exit_code: status.code(),
stdout_tail: stdout_tail.clone(),
stderr_tail: stderr_tail.clone(),
commit_hash: revision_to_history_commit_hash(&end_revision),
};
agent.record_acceptance_attempt(change_id, attempt.clone());
acceptance_history.lock().await.record(change_id, attempt);
acceptance_tail_injected.lock().await.remove(change_id);
let blocker =
crate::events::StalledBlocker::permission_denial("acceptance", denial);
if let Some(ref tx) = event_tx {
let _ = tx
.send(ParallelEvent::ExecutionBlocked {
change_id: change_id.to_string(),
blocker: blocker.clone(),
})
.await;
let _ = tx
.send(ParallelEvent::AcceptanceCompleted {
change_id: change_id.to_string(),
})
.await;
}
return Ok((
crate::orchestration::AcceptanceResult::PermissionStalled { blocker },
attempt_number,
));
}
}
let error_msg = format!(
"Acceptance command failed with exit code: {:?}",
status.code()
);
if let Some(ref tx) = event_tx {
let _ = tx
.send(ParallelEvent::Log(
crate::events::LogEntry::error(&error_msg)
.with_change_id(change_id)
.with_operation("acceptance"),
))
.await;
let _ = tx
.send(ParallelEvent::AcceptanceCompleted {
change_id: change_id.to_string(),
})
.await;
}
return Ok((
crate::orchestration::AcceptanceResult::CommandFailed {
diagnostic: crate::orchestration::acceptance::AcceptanceCommandDiagnostic {
error: error_msg.clone(),
exit_code: status.code(),
stdout_tail: stdout_tail.clone(),
stderr_tail: stderr_tail.clone(),
},
error: error_msg,
findings: tail_findings,
},
0,
));
}
match parse_result {
ParseResult::Pass => {
info!("Acceptance passed for: {}", change_id);
let attempt_number = agent.next_acceptance_attempt_number(change_id);
let attempt = crate::history::AcceptanceAttempt {
attempt: attempt_number,
passed: true,
duration: start_time.elapsed(),
findings: None,
exit_code: status.code(),
stdout_tail: stdout_tail.clone(),
stderr_tail: stderr_tail.clone(),
commit_hash: revision_to_history_commit_hash(&end_revision),
};
agent.record_acceptance_attempt(change_id, attempt.clone());
agent.clear_acceptance_follow_up(change_id);
let mut shared_history = acceptance_history.lock().await;
shared_history.record(change_id, attempt);
shared_history.clear_follow_up_findings(change_id);
drop(shared_history);
acceptance_tail_injected.lock().await.remove(change_id);
if let Some(ref tx) = event_tx {
let _ = tx
.send(ParallelEvent::Log(
crate::events::LogEntry::info("Acceptance test passed")
.with_change_id(change_id)
.with_operation("acceptance")
.with_iteration(attempt_number),
))
.await;
let _ = tx
.send(ParallelEvent::AcceptanceCompleted {
change_id: change_id.to_string(),
})
.await;
}
Ok((crate::orchestration::AcceptanceResult::Pass, attempt_number))
}
ParseResult::Continue => {
info!("Acceptance requires continuation for: {}", change_id);
let attempt_number = agent.next_acceptance_attempt_number(change_id);
let attempt = crate::history::AcceptanceAttempt {
attempt: attempt_number,
passed: false,
duration: start_time.elapsed(),
findings: Some(crate::acceptance::legacy_findings([
"Investigation incomplete - continue later",
])),
exit_code: status.code(),
stdout_tail: stdout_tail.clone(),
stderr_tail: stderr_tail.clone(),
commit_hash: revision_to_history_commit_hash(&end_revision),
};
agent.record_acceptance_attempt(change_id, attempt.clone());
acceptance_history.lock().await.record(change_id, attempt);
acceptance_tail_injected.lock().await.remove(change_id);
if let Some(ref tx) = event_tx {
let _ = tx
.send(ParallelEvent::Log(
crate::events::LogEntry::info("Acceptance test requires continuation")
.with_change_id(change_id)
.with_operation("acceptance")
.with_iteration(attempt_number),
))
.await;
let _ = tx
.send(ParallelEvent::AcceptanceCompleted {
change_id: change_id.to_string(),
})
.await;
}
Ok((
crate::orchestration::AcceptanceResult::Continue,
attempt_number,
))
}
ParseResult::MissingVerdict => {
warn!(
"Acceptance completed without a canonical verdict for: {} (missing-verdict protocol failure)",
change_id
);
let mut evidence =
vec![crate::orchestration::acceptance::MISSING_VERDICT_DIAGNOSTIC.to_string()];
evidence.extend(tail_findings.iter().cloned());
let attempt_number = agent.next_acceptance_attempt_number(change_id);
let attempt = crate::history::AcceptanceAttempt {
attempt: attempt_number,
passed: false,
duration: start_time.elapsed(),
findings: Some(crate::acceptance::legacy_findings(evidence.clone())),
exit_code: status.code(),
stdout_tail: stdout_tail.clone(),
stderr_tail: stderr_tail.clone(),
commit_hash: revision_to_history_commit_hash(&end_revision),
};
agent.record_acceptance_attempt(change_id, attempt.clone());
acceptance_history.lock().await.record(change_id, attempt);
acceptance_tail_injected.lock().await.remove(change_id);
if let Some(ref tx) = event_tx {
let _ = tx
.send(ParallelEvent::Log(
crate::events::LogEntry::warn(
"Acceptance completed without a canonical verdict (missing-verdict \
protocol failure); status-only or waiting output is not a verdict. \
The acceptance agent must wait for owned verification results and \
emit exactly one canonical verdict before exiting.",
)
.with_change_id(change_id)
.with_operation("acceptance")
.with_iteration(attempt_number),
))
.await;
let _ = tx
.send(ParallelEvent::AcceptanceCompleted {
change_id: change_id.to_string(),
})
.await;
}
Ok((
crate::orchestration::AcceptanceResult::MissingVerdict {
findings: tail_findings,
},
attempt_number,
))
}
ParseResult::Stalled { blocker } => {
info!(
"Acceptance reported a validated external blocker for: {} (category {})",
change_id, blocker.category
);
let attempt_number = agent.next_acceptance_attempt_number(change_id);
let mut findings = vec![format!(
"Validated external blocker (category {}): {}",
blocker.category, blocker.next_action
)];
findings.extend(blocker.evidence.iter().cloned());
let attempt = crate::history::AcceptanceAttempt {
attempt: attempt_number,
passed: false,
duration: start_time.elapsed(),
findings: Some(crate::acceptance::legacy_findings(findings)),
exit_code: status.code(),
stdout_tail: stdout_tail.clone(),
stderr_tail: stderr_tail.clone(),
commit_hash: revision_to_history_commit_hash(&end_revision),
};
agent.record_acceptance_attempt(change_id, attempt.clone());
acceptance_history.lock().await.record(change_id, attempt);
acceptance_tail_injected.lock().await.remove(change_id);
if let Some(ref tx) = event_tx {
let _ = tx
.send(ParallelEvent::Log(
crate::events::LogEntry::warn(format!(
"Acceptance stalled ({}) on a validated external blocker: {}",
blocker.category, blocker.next_action
))
.with_change_id(change_id)
.with_operation("acceptance")
.with_iteration(attempt_number),
))
.await;
let _ = tx
.send(ParallelEvent::AcceptanceCompleted {
change_id: change_id.to_string(),
})
.await;
}
Ok((
crate::orchestration::AcceptanceResult::Stalled { blocker },
attempt_number,
))
}
ParseResult::BareBlocker { rejection } => {
warn!(
"Acceptance emitted a bare blocker compatibility token for {}: {}",
change_id,
rejection.reason()
);
let attempt_number = agent.next_acceptance_attempt_number(change_id);
let attempt = crate::history::AcceptanceAttempt {
attempt: attempt_number,
passed: false,
duration: start_time.elapsed(),
findings: Some(crate::acceptance::legacy_findings([
crate::orchestration::acceptance::BARE_BLOCKER_DIAGNOSTIC.to_string(),
rejection.reason(),
])),
exit_code: status.code(),
stdout_tail: stdout_tail.clone(),
stderr_tail: stderr_tail.clone(),
commit_hash: revision_to_history_commit_hash(&end_revision),
};
agent.record_acceptance_attempt(change_id, attempt.clone());
acceptance_history.lock().await.record(change_id, attempt);
acceptance_tail_injected.lock().await.remove(change_id);
if let Some(ref tx) = event_tx {
let _ = tx
.send(ParallelEvent::Log(
crate::events::LogEntry::warn(format!(
"Acceptance emitted a gated compatibility token without a validated \
structured blocker ({}); a stalled hold requires an explicit supported \
category, concrete evidence, next action, and resumability.",
rejection.reason()
))
.with_change_id(change_id)
.with_operation("acceptance")
.with_iteration(attempt_number),
))
.await;
let _ = tx
.send(ParallelEvent::AcceptanceCompleted {
change_id: change_id.to_string(),
})
.await;
}
Ok((
crate::orchestration::AcceptanceResult::BareBlocker { rejection },
attempt_number,
))
}
ParseResult::MalformedFinding { rejection } => {
warn!(
"Acceptance emitted a malformed structured finding for {}: {}",
change_id,
rejection.reason()
);
let attempt_number = agent.next_acceptance_attempt_number(change_id);
let attempt = crate::history::AcceptanceAttempt {
attempt: attempt_number,
passed: false,
duration: start_time.elapsed(),
findings: Some(crate::acceptance::legacy_findings([
crate::orchestration::acceptance::MALFORMED_FINDING_DIAGNOSTIC.to_string(),
rejection.reason(),
])),
exit_code: status.code(),
stdout_tail: stdout_tail.clone(),
stderr_tail: stderr_tail.clone(),
commit_hash: revision_to_history_commit_hash(&end_revision),
};
agent.record_acceptance_attempt(change_id, attempt.clone());
acceptance_history.lock().await.record(change_id, attempt);
acceptance_tail_injected.lock().await.remove(change_id);
if let Some(ref tx) = event_tx {
let _ = tx
.send(ParallelEvent::Log(
crate::events::LogEntry::warn(format!(
"Acceptance emitted a FAIL verdict whose structured finding did not \
validate ({}); runtime will not convert it into a path-only repair \
instruction.",
rejection.reason()
))
.with_change_id(change_id)
.with_operation("acceptance")
.with_iteration(attempt_number),
))
.await;
let _ = tx
.send(ParallelEvent::AcceptanceCompleted {
change_id: change_id.to_string(),
})
.await;
}
Ok((
crate::orchestration::AcceptanceResult::MalformedFinding { rejection },
attempt_number,
))
}
ParseResult::Fail { findings } => {
let findings_for_tasks = if findings.is_empty() {
crate::acceptance::legacy_findings([
crate::orchestration::acceptance::GENERIC_ACCEPTANCE_FAIL_FINDING,
])
} else {
findings
};
let findings_text = crate::acceptance::finding_texts(&findings_for_tasks).join("\n");
let current_denial = crate::permission::classify_permission_denial(&[
stdout_tail.as_deref(),
stderr_tail.as_deref(),
Some(findings_text.as_str()),
]);
if let Some(denial) = ¤t_denial {
let repeated_unresolved = previous_acceptance_denial
.as_ref()
.is_some_and(|previous| previous.signature() == denial.signature())
&& end_revision == start_revision;
if repeated_unresolved {
let attempt_number = agent.next_acceptance_attempt_number(change_id);
let attempt = crate::history::AcceptanceAttempt {
attempt: attempt_number,
passed: false,
duration: start_time.elapsed(),
findings: Some(findings_for_tasks.clone()),
exit_code: status.code(),
stdout_tail: stdout_tail.clone(),
stderr_tail: stderr_tail.clone(),
commit_hash: revision_to_history_commit_hash(&end_revision),
};
agent.record_acceptance_attempt(change_id, attempt.clone());
acceptance_history.lock().await.record(change_id, attempt);
acceptance_tail_injected.lock().await.remove(change_id);
let blocker =
crate::events::StalledBlocker::permission_denial("acceptance", denial);
if let Some(ref tx) = event_tx {
let _ = tx
.send(ParallelEvent::ExecutionBlocked {
change_id: change_id.to_string(),
blocker: blocker.clone(),
})
.await;
let _ = tx
.send(ParallelEvent::AcceptanceCompleted {
change_id: change_id.to_string(),
})
.await;
}
return Ok((
crate::orchestration::AcceptanceResult::PermissionStalled { blocker },
attempt_number,
));
}
}
let blocking_gate_context = findings_for_tasks
.first()
.map(|finding| finding.text().to_string())
.unwrap_or_else(|| "no acceptance findings captured".to_string());
info!(
"Acceptance failed for: {} ({} findings), blocking gate context: {}",
change_id,
findings_for_tasks.len(),
blocking_gate_context
);
let attempt_number = agent.next_acceptance_attempt_number(change_id);
let attempt = crate::history::AcceptanceAttempt {
attempt: attempt_number,
passed: false,
duration: start_time.elapsed(),
findings: Some(findings_for_tasks.clone()),
exit_code: status.code(),
stdout_tail: stdout_tail.clone(),
stderr_tail: stderr_tail.clone(),
commit_hash: revision_to_history_commit_hash(&end_revision),
};
agent.record_acceptance_attempt(change_id, attempt.clone());
let repository_findings =
crate::orchestration::acceptance::repository_findings(&findings_for_tasks);
if !repository_findings.is_empty() {
agent.record_acceptance_follow_up(
change_id,
attempt_number,
repository_findings.clone(),
);
}
let mut shared_history = acceptance_history.lock().await;
shared_history.record(change_id, attempt);
if !repository_findings.is_empty() {
shared_history.set_follow_up_findings(
change_id,
attempt_number,
repository_findings,
);
}
drop(shared_history);
acceptance_tail_injected.lock().await.remove(change_id);
if let Some(ref tx) = event_tx {
let _ = tx
.send(ParallelEvent::Log(
crate::events::LogEntry::warn(format_acceptance_failure_log_message(
&findings_for_tasks,
))
.with_change_id(change_id)
.with_operation("acceptance")
.with_iteration(attempt_number),
))
.await;
let _ = tx
.send(ParallelEvent::AcceptanceCompleted {
change_id: change_id.to_string(),
})
.await;
}
Ok((
crate::orchestration::AcceptanceResult::Fail {
findings: findings_for_tasks,
},
attempt_number,
))
}
}
}
#[cfg(test)]
mod tests {
use super::{
format_acceptance_failure_log_message, mark_acceptance_context_injected,
prepare_acceptance_context_for_apply, resolve_acceptance_state_revision,
run_post_apply_cleanup_review,
};
use crate::agent::AgentRunner;
use crate::ai_command_runner::AiCommandRunner;
use crate::command_queue::CommandQueueConfig;
use crate::config::defaults::default_retry_patterns;
use crate::config::OrchestratorConfig;
use crate::task_parser::TaskProgress;
use std::sync::Arc;
use tempfile::TempDir;
use tokio::process::Command;
use tokio::sync::Mutex;
#[tokio::test]
async fn latest_shared_acceptance_findings_seed_parallel_apply_runner() {
use crate::history::{AcceptanceAttempt, AcceptanceHistory};
use std::time::Duration;
let mut history = AcceptanceHistory::new();
history.record(
"change-a",
AcceptanceAttempt {
attempt: 3,
passed: false,
duration: Duration::from_secs(1),
findings: Some(vec!["fix canonical finding".to_string().into()]),
exit_code: Some(0),
stdout_tail: Some("unstructured noise".to_string()),
stderr_tail: None,
commit_hash: None,
},
);
history.set_follow_up_findings(
"change-a",
3,
vec!["fix canonical finding".to_string().into()],
);
let history = Arc::new(Mutex::new(history));
let injected = Arc::new(Mutex::new(std::collections::HashMap::new()));
let mut first_agent = AgentRunner::new(OrchestratorConfig::default());
prepare_acceptance_context_for_apply(&mut first_agent, "change-a", &history, &injected)
.await;
let context = first_agent.get_acceptance_tail_context_for_apply("change-a");
assert!(context.contains("fix canonical finding"));
assert!(!context.contains("unstructured noise"));
mark_acceptance_context_injected(&first_agent, "change-a", &injected).await;
let mut resumed_agent = AgentRunner::new(OrchestratorConfig::default());
prepare_acceptance_context_for_apply(&mut resumed_agent, "change-a", &history, &injected)
.await;
assert!(resumed_agent
.get_acceptance_tail_context_for_apply("change-a")
.is_empty());
}
#[test]
fn test_progress_commit_message_format() {
let change_id = "add-feature";
let progress = TaskProgress {
completed: 5,
total: 10,
};
let expected = "WIP: add-feature (5/10 tasks)";
let actual = format!(
"WIP: {} ({}/{} tasks)",
change_id, progress.completed, progress.total
);
assert_eq!(actual, expected);
}
#[test]
fn test_progress_commit_message_all_complete() {
let change_id = "fix-bug";
let progress = TaskProgress {
completed: 7,
total: 7,
};
let expected = "WIP: fix-bug (7/7 tasks)";
let actual = format!(
"WIP: {} ({}/{} tasks)",
change_id, progress.completed, progress.total
);
assert_eq!(actual, expected);
}
#[test]
fn test_progress_commit_message_zero_progress() {
let change_id = "new-change";
let progress = TaskProgress {
completed: 0,
total: 5,
};
let expected = "WIP: new-change (0/5 tasks)";
let actual = format!(
"WIP: {} ({}/{} tasks)",
change_id, progress.completed, progress.total
);
assert_eq!(actual, expected);
}
#[test]
fn test_progress_commit_message_special_characters() {
let change_id = "add-web-monitoring-feature";
let progress = TaskProgress {
completed: 50,
total: 70,
};
let expected = "WIP: add-web-monitoring-feature (50/70 tasks)";
let actual = format!(
"WIP: {} ({}/{} tasks)",
change_id, progress.completed, progress.total
);
assert_eq!(actual, expected);
}
#[test]
fn test_format_acceptance_failure_log_message_includes_blocking_gate_context() {
let message = format_acceptance_failure_log_message(&[
"archive-readiness gate failed: cargo clippy -- -D warnings (src/lib.rs:42)"
.to_string()
.into(),
"second finding".to_string().into(),
]);
assert!(
message.contains("Acceptance failed (2 findings), blocking gate context:"),
"message should include finding count and blocking gate context prefix"
);
assert!(
message.contains("cargo clippy -- -D warnings"),
"message should preserve gate-specific failure context"
);
}
#[test]
fn test_format_acceptance_failure_log_message_handles_empty_findings() {
let message = format_acceptance_failure_log_message(&[]);
assert_eq!(
message,
"Acceptance failed (0 findings), blocking gate context: no acceptance findings captured"
);
}
#[test]
fn test_resolve_acceptance_state_revision_prefers_end_revision() {
let resolved = resolve_acceptance_state_revision("start-rev", Some("end-rev".to_string()));
assert_eq!(resolved, "end-rev");
}
#[test]
fn test_resolve_acceptance_state_revision_falls_back_to_start_revision() {
let resolved = resolve_acceptance_state_revision("start-rev", None);
assert_eq!(resolved, "start-rev");
}
#[test]
fn test_progress_check_condition() {
let old_progress = TaskProgress {
completed: 3,
total: 10,
};
let new_progress_same = TaskProgress {
completed: 3,
total: 10,
};
let new_progress_increased = TaskProgress {
completed: 5,
total: 10,
};
let new_progress_decreased = TaskProgress {
completed: 2,
total: 10,
};
assert!(new_progress_same.completed <= old_progress.completed);
assert!(new_progress_increased.completed > old_progress.completed);
assert!(new_progress_decreased.completed <= old_progress.completed);
}
async fn init_test_git_repo(repo_root: &std::path::Path) {
Command::new("git")
.args(["init", "-b", "main"])
.current_dir(repo_root)
.output()
.await
.unwrap();
Command::new("git")
.args(["config", "user.email", "test@example.com"])
.current_dir(repo_root)
.output()
.await
.unwrap();
Command::new("git")
.args(["config", "user.name", "Test User"])
.current_dir(repo_root)
.output()
.await
.unwrap();
std::fs::write(repo_root.join("README.md"), "base\n").unwrap();
Command::new("git")
.args(["add", "README.md"])
.current_dir(repo_root)
.output()
.await
.unwrap();
Command::new("git")
.args(["commit", "-m", "base"])
.current_dir(repo_root)
.output()
.await
.unwrap();
}
fn test_ai_runner() -> AiCommandRunner {
let queue_config = CommandQueueConfig {
acceptance_max_runtime_secs:
crate::config::defaults::DEFAULT_ACCEPTANCE_MAX_RUNTIME_SECS,
stagger_delay_ms: 0,
max_retries: 0,
retry_delay_ms: 0,
retry_error_patterns: default_retry_patterns(),
retry_if_duration_under_secs: 0,
inactivity_timeout_secs: 0,
inactivity_kill_grace_secs: 1,
inactivity_timeout_max_retries: 0,
strict_process_cleanup: true,
max_runtime_secs: 0,
};
let shared_stagger_state = Arc::new(Mutex::new(None));
AiCommandRunner::new(queue_config, shared_stagger_state)
}
#[tokio::test]
async fn test_post_apply_cleanup_review_succeeds_with_single_clean_marker() {
let temp_dir = TempDir::new().unwrap();
init_test_git_repo(temp_dir.path()).await;
std::fs::write(temp_dir.path().join("dirty.txt"), "dirty\n").unwrap();
let config = OrchestratorConfig {
acceptance_command: Some(
"sh -c 'git add dirty.txt && git commit -m cleanup && echo CLEANUP_REVIEW: CLEAN'"
.to_string(),
),
..Default::default()
};
let ai_runner = test_ai_runner();
run_post_apply_cleanup_review("change-a", temp_dir.path(), &config, &ai_runner, None, None)
.await
.expect("cleanup review should succeed");
let (is_dirty, status) =
crate::vcs::git::commands::has_uncommitted_changes(temp_dir.path())
.await
.unwrap();
assert!(
!is_dirty,
"worktree must be clean after successful cleanup review: {status}"
);
}
#[tokio::test]
async fn test_post_apply_cleanup_review_fails_when_marker_missing() {
let temp_dir = TempDir::new().unwrap();
init_test_git_repo(temp_dir.path()).await;
std::fs::write(temp_dir.path().join("dirty.txt"), "dirty\n").unwrap();
let config = OrchestratorConfig {
acceptance_command: Some(
"sh -c 'git add dirty.txt; git diff --cached --quiet || git commit -m cleanup; \
echo done'"
.to_string(),
),
..Default::default()
};
let ai_runner = test_ai_runner();
let err = run_post_apply_cleanup_review(
"change-a",
temp_dir.path(),
&config,
&ai_runner,
None,
None,
)
.await
.expect_err("cleanup review must fail without marker");
let message = err.to_string();
assert!(
message.contains("marker_missing"),
"error should name the missing-marker failure kind: {message}"
);
assert!(
message.contains("3 operation attempts"),
"error should report the exhausted attempt count: {message}"
);
}
mod cleanup_review_recovery {
use super::super::run_post_apply_cleanup_review;
use crate::ai_command_runner::AiCommandRunner;
use crate::command_queue::CommandQueueConfig;
use crate::config::OrchestratorConfig;
use crate::error::OrchestratorError;
use std::path::Path;
use std::sync::Arc;
use tempfile::TempDir;
use tokio::sync::Mutex;
use tokio_util::sync::CancellationToken;
fn ai_runner() -> AiCommandRunner {
let queue_config = CommandQueueConfig {
acceptance_max_runtime_secs:
crate::config::defaults::DEFAULT_ACCEPTANCE_MAX_RUNTIME_SECS,
stagger_delay_ms: 0,
max_retries: 0,
retry_delay_ms: 0,
retry_error_patterns: Vec::new(),
retry_if_duration_under_secs: 0,
inactivity_timeout_secs: 0,
inactivity_kill_grace_secs: 1,
inactivity_timeout_max_retries: 0,
strict_process_cleanup: true,
max_runtime_secs: 0,
};
AiCommandRunner::new(queue_config, Arc::new(Mutex::new(None)))
}
fn git(repo: &Path, args: &[&str]) {
let output = std::process::Command::new("git")
.args(args)
.current_dir(repo)
.output()
.expect("git should run");
assert!(
output.status.success(),
"git {args:?} failed: {}",
String::from_utf8_lossy(&output.stderr)
);
}
fn dirty_worktree() -> TempDir {
let temp_dir = TempDir::new().unwrap();
let repo = temp_dir.path();
git(repo, &["init", "-b", "main"]);
git(repo, &["config", "user.email", "test@example.com"]);
git(repo, &["config", "user.name", "Test User"]);
std::fs::write(repo.join("README.md"), "base\n").unwrap();
git(repo, &["add", "README.md"]);
git(repo, &["commit", "-m", "base"]);
std::fs::write(repo.join("leftover.txt"), "apply artifact\n").unwrap();
temp_dir
}
fn fixture(state: &Path, body: &str) -> String {
let script = state.join("cleanup-review.sh");
std::fs::write(
&script,
format!(
"#!/bin/sh\n\
STATE={state}\n\
ATTEMPT=$(cat \"$STATE/attempts\" 2>/dev/null || echo 0)\n\
ATTEMPT=$((ATTEMPT+1))\n\
echo $ATTEMPT > \"$STATE/attempts\"\n\
printf '%s' \"$1\" > \"$STATE/prompt-$ATTEMPT.txt\"\n\
{body}\n",
state = state.display(),
body = body
),
)
.unwrap();
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
std::fs::set_permissions(&script, std::fs::Permissions::from_mode(0o755)).unwrap();
}
format!("{} {{prompt}}", script.display())
}
fn config(command: String) -> OrchestratorConfig {
OrchestratorConfig {
acceptance_command: Some(command),
..Default::default()
}
}
fn attempts(state: &Path) -> u32 {
std::fs::read_to_string(state.join("attempts"))
.map(|text| text.trim().parse().unwrap_or(0))
.unwrap_or(0)
}
fn prompt(state: &Path, attempt: u32) -> String {
std::fs::read_to_string(state.join(format!("prompt-{attempt}.txt")))
.unwrap_or_else(|e| panic!("attempt {attempt} prompt should exist: {e}"))
}
async fn is_dirty(repo: &Path) -> bool {
crate::vcs::git::commands::has_uncommitted_changes(repo)
.await
.unwrap()
.0
}
#[cfg_attr(windows, ignore)]
#[tokio::test]
async fn a_corrective_attempt_recovers_a_dirty_worktree() {
let repo_dir = dirty_worktree();
let state = TempDir::new().unwrap();
let command = fixture(
state.path(),
"if [ \"$ATTEMPT\" = \"1\" ]; then echo 'reviewed, nothing done'; exit 0; fi\n\
git add leftover.txt && git commit -q -m cleanup && echo 'CLEANUP_REVIEW: CLEAN'",
);
run_post_apply_cleanup_review(
"change-a",
repo_dir.path(),
&config(command),
&ai_runner(),
None,
None,
)
.await
.expect("the corrective attempt must recover the handoff");
assert_eq!(
attempts(state.path()),
2,
"exactly one correction was needed"
);
assert!(!is_dirty(repo_dir.path()).await);
let first = prompt(state.path(), 1);
assert!(
!first.contains("<cleanup_review_correction>"),
"the initial attempt has nothing to correct"
);
let second = prompt(state.path(), 2);
assert!(second.contains("<cleanup_review_correction>"), "{second}");
assert!(
second.contains("\"failure_kind\":\"marker_missing\""),
"{second}"
);
assert!(
second.contains("\"standalone_clean_marker_count\":0"),
"{second}"
);
assert!(
second.contains("leftover.txt"),
"the corrective prompt must carry fresh porcelain evidence: {second}"
);
assert!(second.contains("Success is decided by Conflux"), "{second}");
}
#[cfg_attr(windows, ignore)]
#[tokio::test]
async fn a_marker_without_a_clean_repository_is_not_success() {
let repo_dir = dirty_worktree();
let state = TempDir::new().unwrap();
let command = fixture(state.path(), "echo 'CLEANUP_REVIEW: CLEAN'");
let error = run_post_apply_cleanup_review(
"change-a",
repo_dir.path(),
&config(command),
&ai_runner(),
None,
None,
)
.await
.expect_err("a marker can never override Git state");
let message = error.to_string();
assert!(message.contains("dirty_remains"), "{message}");
assert!(message.contains("leftover.txt"), "{message}");
assert_eq!(attempts(state.path()), 3);
assert!(
prompt(state.path(), 2).contains("\"failure_kind\":\"dirty_remains\""),
"the correction must name the observed failure"
);
assert!(
is_dirty(repo_dir.path()).await,
"exhaustion preserves the managed workspace for explicit retry"
);
}
#[cfg_attr(windows, ignore)]
#[tokio::test]
async fn a_clean_repository_without_exactly_one_marker_is_not_success() {
let repo_dir = dirty_worktree();
let state = TempDir::new().unwrap();
let command = fixture(
state.path(),
"git add leftover.txt >/dev/null 2>&1; \
git diff --cached --quiet || git commit -q -m cleanup; \
echo 'CLEANUP_REVIEW: CLEAN'; echo 'CLEANUP_REVIEW: CLEAN'",
);
let error = run_post_apply_cleanup_review(
"change-a",
repo_dir.path(),
&config(command),
&ai_runner(),
None,
None,
)
.await
.expect_err("a clean worktree never invents the marker contract");
let message = error.to_string();
assert!(message.contains("marker_duplicate"), "{message}");
assert!(
message.contains("standalone_clean_marker_count: 2"),
"{message}"
);
assert_eq!(
attempts(state.path()),
3,
"exactly three operation attempts"
);
}
#[cfg_attr(windows, ignore)]
#[tokio::test]
async fn exhaustion_reports_the_attempt_count_and_latest_diagnosis() {
let repo_dir = dirty_worktree();
let state = TempDir::new().unwrap();
let command = fixture(
state.path(),
"echo \"cleanup crashed on attempt $ATTEMPT\" >&2; exit 5",
);
let error = run_post_apply_cleanup_review(
"change-a",
repo_dir.path(),
&config(command),
&ai_runner(),
None,
None,
)
.await
.expect_err("three failed operation attempts are terminal");
let message = error.to_string();
assert!(message.contains("3 operation attempts"), "{message}");
assert!(
message.contains("failure_kind: command_failed"),
"{message}"
);
assert!(message.contains("exit_code: 5"), "{message}");
assert!(
message.contains("cleanup crashed on attempt 3"),
"the terminal error reports the latest diagnosis: {message}"
);
assert_eq!(
attempts(state.path()),
super::super::MAX_CLEANUP_REVIEW_RETRIES + 1,
"no fourth operation attempt may start"
);
assert!(is_dirty(repo_dir.path()).await);
}
#[tokio::test]
async fn an_already_cancelled_token_starts_no_attempt() {
let repo_dir = dirty_worktree();
let state = TempDir::new().unwrap();
let command = fixture(state.path(), "echo 'CLEANUP_REVIEW: CLEAN'");
let cancel = CancellationToken::new();
cancel.cancel();
let error = run_post_apply_cleanup_review(
"change-a",
repo_dir.path(),
&config(command),
&ai_runner(),
Some(&cancel),
None,
)
.await
.expect_err("cancellation is an intentional stop, not a cleanup success");
assert!(
error.is_cancellation(),
"the run boundary must be able to classify the stop without parsing text: {error}"
);
assert!(
matches!(
&error,
OrchestratorError::Cancelled { operation, change_id, .. }
if operation == "cleanup-review" && change_id == "change-a"
),
"the typed stop names the operation and change: {error}"
);
assert!(
error.to_string().contains("Cancelled cleanup-review"),
"the existing rendering is unchanged: {error}"
);
assert_eq!(
attempts(state.path()),
0,
"explicit cancellation must not start an attempt"
);
}
#[cfg_attr(windows, ignore)]
#[cfg_attr(not(feature = "heavy-tests"), ignore)]
#[tokio::test]
async fn cancellation_terminates_the_child_and_starts_no_further_attempt() {
let repo_dir = dirty_worktree();
let state = TempDir::new().unwrap();
let command = fixture(
state.path(),
"touch \"$STATE/started\"; sleep 30; echo 'CLEANUP_REVIEW: CLEAN'",
);
let cancel = CancellationToken::new();
let started_marker = state.path().join("started");
let waiter = cancel.clone();
let starter = tokio::spawn(async move {
for _ in 0..12_000 {
if started_marker.exists() {
break;
}
tokio::time::sleep(std::time::Duration::from_millis(10)).await;
}
waiter.cancel();
});
let error = run_post_apply_cleanup_review(
"change-a",
repo_dir.path(),
&config(command),
&ai_runner(),
Some(&cancel),
None,
)
.await
.expect_err("cancellation is an intentional stop, not a cleanup success");
starter.abort();
assert!(
error.is_cancellation(),
"cancellation keeps the typed intentional-stop routing: {error}"
);
assert_eq!(
attempts(state.path()),
1,
"explicit cancellation must not start a corrective attempt"
);
}
#[cfg_attr(windows, ignore)]
#[tokio::test]
async fn a_classified_permission_denial_holds_without_consuming_the_failure_budget() {
let repo_dir = dirty_worktree();
let state = TempDir::new().unwrap();
let command = fixture(
state.path(),
"echo 'Tool access denied: Bash(git commit)'; exit 0",
);
let error = run_post_apply_cleanup_review(
"change-a",
repo_dir.path(),
&config(command),
&ai_runner(),
None,
None,
)
.await
.expect_err("a classified denial enters the non-terminal hold");
assert!(
matches!(error, OrchestratorError::PermissionStalled { .. }),
"cleanup permission denial reuses the existing non-terminal hold, got {error:?}"
);
assert_eq!(
attempts(state.path()),
1,
"permission denial starts no corrective attempt and consumes no generic budget"
);
assert!(
is_dirty(repo_dir.path()).await,
"the managed workspace is preserved for an explicit retry"
);
}
#[cfg_attr(windows, ignore)]
#[tokio::test]
async fn a_restarted_run_derives_cleanup_from_workspace_evidence_with_a_fresh_budget() {
let repo_dir = dirty_worktree();
let first_state = TempDir::new().unwrap();
let first = fixture(first_state.path(), "echo nope; exit 1");
run_post_apply_cleanup_review(
"change-a",
repo_dir.path(),
&config(first),
&ai_runner(),
None,
None,
)
.await
.expect_err("the first run exhausts its budget");
let workspace_entries: Vec<String> = std::fs::read_dir(repo_dir.path())
.unwrap()
.map(|entry| entry.unwrap().file_name().to_string_lossy().to_string())
.filter(|name| name != ".git")
.collect();
assert_eq!(
workspace_entries
.iter()
.filter(|name| name.as_str() != "README.md" && name.as_str() != "leftover.txt")
.count(),
0,
"no durable retry artifact may be created: {workspace_entries:?}"
);
let second_state = TempDir::new().unwrap();
let second = fixture(
second_state.path(),
"git add leftover.txt && git commit -q -m cleanup && echo 'CLEANUP_REVIEW: CLEAN'",
);
run_post_apply_cleanup_review(
"change-a",
repo_dir.path(),
&config(second),
&ai_runner(),
None,
None,
)
.await
.expect("a restarted run gets a fresh active-run cleanup budget");
assert_eq!(
attempts(second_state.path()),
1,
"the restarted run starts from a clean budget"
);
assert!(
!prompt(second_state.path(), 1).contains("<cleanup_review_correction>"),
"no prior-run diagnosis may be restored from outside the workspace"
);
}
#[cfg_attr(windows, ignore)]
#[tokio::test]
async fn a_command_failure_that_also_breaks_status_reports_both() {
let repo_dir = dirty_worktree();
let state = TempDir::new().unwrap();
let command = fixture(
state.path(),
"rm -rf .git; echo 'cleanup crashed' >&2; exit 5",
);
let error = run_post_apply_cleanup_review(
"change-a",
repo_dir.path(),
&config(command),
&ai_runner(),
None,
None,
)
.await
.expect_err("a broken repository is never a handoff-ready worktree");
let message = error.to_string();
assert!(
message.contains("failure_kind: command_failed"),
"the primary failure kind is preserved: {message}"
);
assert!(message.contains("exit_code: 5"), "{message}");
assert!(
message.contains("status_error: "),
"the simultaneous status-inspection failure must be reported: {message}"
);
assert!(message.contains("status inspection failed"), "{message}");
assert!(
!message.contains("| status: "),
"an unanswerable query must not present itself as observed status: {message}"
);
let second = prompt(state.path(), 2);
assert!(
second.contains("\"failure_kind\":\"command_failed\""),
"{second}"
);
assert!(
second.contains("\"status_inspection_error\""),
"the corrective prompt must say cleanliness is unproven: {second}"
);
assert!(
!second.contains("\"current_porcelain_status\""),
"no status may be claimed when the query failed: {second}"
);
}
#[cfg_attr(windows, ignore)]
#[tokio::test]
async fn a_marker_failure_that_also_breaks_status_reports_both() {
let repo_dir = dirty_worktree();
let state = TempDir::new().unwrap();
let command = fixture(state.path(), "rm -rf .git; echo 'reviewed'; exit 0");
let error = run_post_apply_cleanup_review(
"change-a",
repo_dir.path(),
&config(command),
&ai_runner(),
None,
None,
)
.await
.expect_err("a missing marker with unprovable cleanliness is terminal");
let message = error.to_string();
assert!(
message.contains("failure_kind: marker_missing"),
"the marker contract still owns the primary kind: {message}"
);
assert!(
message.contains("standalone_clean_marker_count: 0"),
"{message}"
);
assert!(
message.contains("status_error: "),
"the simultaneous status-inspection failure must be reported: {message}"
);
let second = prompt(state.path(), 2);
assert!(
second.contains("\"failure_kind\":\"marker_missing\""),
"{second}"
);
assert!(second.contains("\"status_inspection_error\""), "{second}");
}
#[test]
fn the_marker_scanner_state_is_fixed_size_regardless_of_stream_length() {
use crate::agent::CleanupMarkerScanner;
let mut scanner = CleanupMarkerScanner::new();
for i in 0..200_000 {
scanner.observe(&format!("noise line {i} with some padding text"));
}
scanner.observe("CLEANUP_REVIEW: CLEAN");
for i in 0..200_000 {
scanner.observe(&format!("trailing noise {i}"));
}
assert_eq!(scanner.count(), 1, "the marker contract still holds");
assert!(
std::mem::size_of::<CleanupMarkerScanner>() <= 2 * std::mem::size_of::<usize>(),
"marker state must stay bounded, got {} bytes",
std::mem::size_of::<CleanupMarkerScanner>()
);
}
#[test]
fn the_marker_scanner_ignores_fenced_markers_across_chunk_boundaries() {
use crate::agent::CleanupMarkerScanner;
let mut scanner = CleanupMarkerScanner::new();
for chunk in [
"```",
"CLEANUP_REVIEW: CLEAN",
"```",
" CLEANUP_REVIEW: CLEAN ",
"prefix CLEANUP_REVIEW: CLEAN",
] {
scanner.observe(chunk);
}
assert_eq!(
scanner.count(),
1,
"only the standalone unfenced marker counts, even when the fence spans chunks"
);
}
#[cfg_attr(windows, ignore)]
#[tokio::test]
async fn a_large_stream_still_validates_the_marker_exactly_once() {
let repo_dir = dirty_worktree();
let state = TempDir::new().unwrap();
let command = fixture(
state.path(),
"echo '```'; echo 'CLEANUP_REVIEW: CLEAN'; echo '```'; \
i=0; while [ $i -lt 4000 ]; do echo \"filler line $i\"; i=$((i+1)); done; \
git add leftover.txt && git commit -q -m cleanup && echo 'CLEANUP_REVIEW: CLEAN'",
);
run_post_apply_cleanup_review(
"change-a",
repo_dir.path(),
&config(command),
&ai_runner(),
None,
None,
)
.await
.expect("one standalone marker outside fences plus a clean worktree is success");
assert_eq!(attempts(state.path()), 1, "no correction was needed");
assert!(!is_dirty(repo_dir.path()).await);
}
}
}