use std::fmt::Write as _;
use std::path::Path;
use std::sync::Arc;
use crate::agent::role::DIAGNOSTICS_ROLE;
use crate::git::commands::{
has_unstaged_changes, run_git_add_all, run_git_head, run_git_status, run_git_worktree_snapshot,
run_git_write_tree,
};
use crate::pipeline::board::Ticket;
use crate::prompt::{load_prompt, load_prompt_sections, substitute};
use crate::tools::shell::{ShellMode, ShellTool};
use crate::{DiagnosticsCommands, Role, Workspace};
use super::{
AgentSlot, ExtractionMode, FinalizeOutcome, ParallelVerdict, TicketPhase, TransitionCtx,
agent_slot_from_roster_row, bounce_to_development, build_agent_slots, build_round_grouping,
comment_and_transition, debug, guard_job_phase, info, insert_round_slots, pause_freezing,
read_roster_or_bail, render_joint_comment, reset_phase_attempt, run_parallel_agents,
sync_phase_job_task, warn,
};
pub(crate) const STAGE: &str = "Verification";
const ACTIVE_PHASE: TicketPhase = TicketPhase::Verification;
const SUCCESS_PHASE: TicketPhase = TicketPhase::InSanitation;
const VERIFICATION_THRESHOLD: u8 = 9;
const TESTER_COUNT: usize = 1;
const SKIPPED_REVIEW_NOTE: &str = "\
**Code review skipped** — the content is identical to what the reviewers were \
last handed on this ticket (same commit, same working tree, no changes since), \
so they were not dispatched again. The functional check still ran.\n\n";
#[derive(Copy, Clone)]
struct Participant {
role: Role,
kind: crate::jobs::AgentKind,
prompt_template: &'static str,
angles_path: &'static str,
extraction_prompt_path: &'static str,
}
const REVIEWER: Participant = Participant {
role: Role::Reviewer,
kind: crate::jobs::AgentKind::Reviewer,
prompt_template: "review.md",
angles_path: "review_angles.md",
extraction_prompt_path: "extraction/reviewer.md",
};
const TESTER: Participant = Participant {
role: Role::Qa,
kind: crate::jobs::AgentKind::Tester,
prompt_template: "qa.md",
angles_path: "qa_angles.md",
extraction_prompt_path: "extraction/qa.md",
};
impl Participant {
fn prompt(self, ticket: &Ticket, commands_brief: &str) -> String {
substitute(
&load_prompt(self.prompt_template),
&[
("{{agent_response}}", engineer_response(ticket)),
("{{project_commands}}", commands_brief),
],
)
}
fn angles(self) -> Vec<String> {
load_prompt_sections(self.angles_path)
}
fn extraction_prompt(self) -> String {
load_prompt(self.extraction_prompt_path)
}
}
fn engineer_response(ticket: &Ticket) -> &str {
ticket
.comments
.iter()
.rev()
.find(|c| c.role == Role::Engineer.as_str())
.map_or("(no output)", |c| c.content.as_str())
}
pub(crate) async fn run(ticket: Arc<Ticket>, ws: Workspace, job_id: String) {
if guard_job_phase(&ticket.id, ACTIVE_PHASE, &job_id).await {
return;
}
let commands = match crate::workspace::store().get_diagnostics(&ws.name).await {
Ok(Some(cmds)) if !cmds.is_empty() => Some(cmds),
Ok(_) => None,
Err(e) => {
warn!(
ticket = %ticket.id,
error = %e,
"Failed to load the workspace's project commands — resetting for a fresh attempt",
);
reset_phase_attempt(
&ticket,
ACTIVE_PHASE,
&job_id,
"project commands load failure",
&format!(
"Could not load the workspace's project commands due to a database error: {e}"
),
)
.await;
return;
}
};
let Some(roster) = read_roster_or_bail(&ticket.id, &job_id).await else {
return;
};
let resume = !roster.is_empty();
let plan = if resume {
let Some(plan) = RoundPlan::from_roster(&roster) else {
warn!(
ticket = %ticket.id,
job = %job_id,
rows = roster.len(),
kinds = %roster.iter().map(|r| r.kind.as_str()).collect::<Vec<_>>().join(","),
"Unreadable verification roster — re-driving the round from scratch",
);
let _ = crate::jobs::terminalize_job(&crate::session::store().conn, &job_id).await;
return;
};
plan
} else {
let Some(plan) = RoundPlan::fresh(&ticket, &ws, &job_id, commands.as_ref()).await else {
return;
};
plan
};
execute_round(ticket, ws, job_id, plan, resume, commands).await;
}
struct RoundPlan {
reviewers: Vec<AgentSlot>,
testers: Vec<AgentSlot>,
}
impl RoundPlan {
async fn fresh(
ticket: &Ticket,
ws: &Workspace,
job_id: &str,
commands: Option<&DiagnosticsCommands>,
) -> Option<Self> {
let commands_brief = commands_brief(commands);
let tester_prompt = TESTER.prompt(ticket, &commands_brief);
let (reviewer_prompt, reviewers) = if skip_review_for_identical_content(ticket, ws).await {
info!(
ticket = %ticket.id,
"Content identical to the reviewed base — dispatching the functional tester only",
);
(None, Vec::new())
} else {
let prompt = REVIEWER.prompt(ticket, &commands_brief);
let slots = build_agent_slots(
&ticket.id,
REVIEWER.role,
&prompt,
&REVIEWER.angles(),
0,
compute_reviewer_count(ticket, ws.as_path()).await,
);
(Some(prompt), slots)
};
let testers = build_agent_slots(
&ticket.id,
TESTER.role,
&tester_prompt,
&TESTER.angles(),
i64::try_from(reviewers.len()).unwrap_or(i64::MAX),
TESTER_COUNT,
);
let task = reviewer_prompt.as_deref().unwrap_or(&tester_prompt);
sync_phase_job_task(&crate::session::store().conn, job_id, task).await;
let cohorts = [
(REVIEWER.kind, reviewers.as_slice()),
(TESTER.kind, testers.as_slice()),
];
if let Err(e) = insert_round_slots(job_id, &cohorts).await {
warn!(
ticket = %ticket.id,
job = %job_id,
error = %e,
"Failed to write the verification roster — round not started",
);
return None;
}
Some(Self { reviewers, testers })
}
fn from_roster(roster: &[crate::jobs::AgentRow]) -> Option<Self> {
let cohort = |kind: crate::jobs::AgentKind| -> Vec<AgentSlot> {
roster
.iter()
.filter(|row| row.kind == kind.as_str())
.map(agent_slot_from_roster_row)
.collect()
};
let reviewers = cohort(REVIEWER.kind);
let testers = cohort(TESTER.kind);
if testers.len() != TESTER_COUNT || reviewers.len() + testers.len() != roster.len() {
return None;
}
Some(Self { reviewers, testers })
}
}
struct ReviewBase {
head: String,
tree: String,
}
fn commands_brief(commands: Option<&DiagnosticsCommands>) -> String {
let Some(commands) = commands else {
return load_prompt("pipeline/round_commands_none.md");
};
let list = commands
.commands()
.iter()
.filter_map(|(label, cmd)| cmd.map(|cmd| format!("- {label}: `{cmd}`")))
.collect::<Vec<_>>()
.join("\n");
substitute(
&load_prompt("pipeline/round_commands.md"),
&[("{{commands}}", &list)],
)
}
struct ProjectCommands {
section: String,
passed: bool,
}
async fn run_project_commands(
commands: Option<&DiagnosticsCommands>,
ws: &Workspace,
) -> ProjectCommands {
let Some(commands) = commands else {
return ProjectCommands {
section: format!(
"**Project commands**\n\n{}",
load_prompt("pipeline/commands_none.md")
),
passed: true,
};
};
let mut lines = String::new();
let mut failed_at: Option<&str> = None;
for (label, cmd) in commands.commands() {
let Some(cmd) = cmd else {
continue;
};
let started = std::time::Instant::now();
let elapsed = || started.elapsed().as_secs_f64();
match ShellTool::new(ShellMode::Full)
.execute_with_status(ws, serde_json::json!({ "command": cmd }))
.await
{
Ok((_output, Some(0))) => {
let _ = writeln!(lines, "- {label} (`{cmd}`): PASSED in {:.1}s", elapsed());
}
Ok((output, _exit_code)) => {
let display = if output.is_empty() {
"(no output)".to_string()
} else {
output
};
let _ = writeln!(
lines,
"- {label} (`{cmd}`): FAILED in {:.1}s\n\n```\n{display}\n```",
elapsed(),
);
failed_at = Some(label);
break;
}
Err(e) => {
let _ = writeln!(
lines,
"- {label} (`{cmd}`): FAILED in {:.1}s\n\n```\n{e}\n```",
elapsed(),
);
failed_at = Some(label);
break;
}
}
}
crate::tools::shell::cleanup_agent_spills(DIAGNOSTICS_ROLE);
let footer = match failed_at {
Some(label) => format!("{} `{label}`", load_prompt("pipeline/commands_failed.md")),
None => load_prompt("pipeline/commands_passed.md"),
};
ProjectCommands {
section: format!("**Project commands**\n\n{}\n\n{footer}", lines.trim_end()),
passed: failed_at.is_none(),
}
}
async fn execute_round(
ticket: Arc<Ticket>,
ws: Workspace,
job_id: String,
plan: RoundPlan,
resume: bool,
commands: Option<DiagnosticsCommands>,
) {
if resume {
let not_done: Vec<String> = plan
.reviewers
.iter()
.chain(&plan.testers)
.filter(|s| s.status != crate::jobs::RowStatus::Done)
.map(|s| s.agent_id.clone())
.collect();
if let Err(e) =
crate::jobs::rearm_roster_launched(&crate::session::store().conn, &job_id, ¬_done)
.await
{
warn!(ticket = %ticket.id, job = %job_id, error = %e, "Failed to re-arm resumed verification slots");
}
}
let reviewer_count = plan.reviewers.len();
let tester_count = plan.testers.len();
let command_count = commands.as_ref().map_or(0, |c| {
c.commands().iter().filter(|(_, cmd)| cmd.is_some()).count()
});
info!(
ticket = %ticket.id,
reviewers = reviewer_count,
testers = tester_count,
commands = command_count,
"Dispatching {reviewer_count} reviewer(s), {tester_count} tester(s) and {command_count} project command(s) in parallel",
);
let review_skipped = plan.reviewers.is_empty();
let base = if review_skipped {
None
} else {
match crate::git::commands::run_git_worktree_identity(ws.as_path()).await {
Ok((head, tree)) => Some(ReviewBase { head, tree }),
Err(e) => {
debug!(error = %e, "Could not read the reviewed content identity");
None
}
}
};
let ((reviewer_results, reviewer_paused), (tester_results, tester_paused), commands_outcome) = tokio::join!(
dispatch_cohort(&ticket, &ws, REVIEWER, &plan.reviewers, &job_id, resume),
dispatch_cohort(&ticket, &ws, TESTER, &plan.testers, &job_id, resume),
run_project_commands(commands.as_ref(), &ws),
);
if guard_job_phase(&ticket.id, ACTIVE_PHASE, &job_id).await {
return;
}
if reviewer_paused || tester_paused {
pause_freezing(&ticket, &job_id).await;
return;
}
finalize_round(
&ws,
&ticket,
&reviewer_results,
&tester_results,
base,
review_skipped,
&job_id,
&commands_outcome,
)
.await;
}
async fn dispatch_cohort(
ticket: &Arc<Ticket>,
ws: &Workspace,
participant: Participant,
slots: &[AgentSlot],
job_id: &str,
resume: bool,
) -> (Vec<ParallelVerdict>, bool) {
if slots.is_empty() {
return (Vec::new(), false);
}
run_parallel_agents(
ticket,
ws,
participant.role,
&participant.extraction_prompt(),
ExtractionMode::ScoreVerdict,
job_id,
slots,
ACTIVE_PHASE,
resume,
)
.await
}
#[expect(clippy::too_many_arguments)]
async fn finalize_round(
ws: &Workspace,
ticket: &Ticket,
reviewer_results: &[ParallelVerdict],
tester_results: &[ParallelVerdict],
base: Option<ReviewBase>,
review_skipped: bool,
job_id: &str,
commands: &ProjectCommands,
) {
let results: Vec<ParallelVerdict> = reviewer_results
.iter()
.chain(tester_results)
.cloned()
.collect();
let transitioned =
process_round_verdicts(ws, ticket, &results, review_skipped, job_id, commands).await;
if !review_skipped {
record_reviewed_base(ws, &ticket.id, base, transitioned).await;
}
}
async fn process_round_verdicts(
ws: &Workspace,
ticket: &Ticket,
results: &[ParallelVerdict],
review_skipped: bool,
job_id: &str,
commands: &ProjectCommands,
) -> bool {
let technical_failure = results.iter().any(ParallelVerdict::is_technical_failure);
let rework_failure = !technical_failure
&& (!commands.passed
|| results
.iter()
.any(|r| matches!(r, ParallelVerdict::Verdict(v) if !verdict_passes(v))));
if crate::shutdown::aborting() {
info!(
ticket = %ticket.id,
stage = STAGE,
"Verification round cut short by drain — job stays launched for boot resume",
);
return false;
}
if technical_failure {
let comment =
format!("{STAGE} could not complete the round (a participant did not respond).");
reset_phase_attempt(
ticket,
ACTIVE_PHASE,
job_id,
"verification failure",
&comment,
)
.await;
return false;
}
let (round, outcome) = build_round_grouping(
STAGE,
results,
Role::Qa,
review_skipped,
ws,
&ticket.id,
&ticket.title,
)
.await;
let mut verdict = render_joint_comment(
&round,
&outcome,
&crate::consensus::ItemTable::new(&round.issues),
);
if review_skipped {
verdict.insert_str(0, SKIPPED_REVIEW_NOTE);
}
let comment = format!("{}\n\n---\n\n{verdict}", commands.section);
if !rework_failure {
return apply_clean_round(ticket, &comment, job_id).await;
}
let outcome = bounce_to_development(
ticket,
ACTIVE_PHASE,
STAGE,
STAGE,
ACTIVE_PHASE.as_ref(),
&comment,
job_id,
)
.await;
matches!(outcome, FinalizeOutcome::Applied)
}
async fn apply_clean_round(ticket: &Ticket, comment: &str, job_id: &str) -> bool {
if !matches!(
comment_and_transition(
TransitionCtx::buffered(
ticket,
ACTIVE_PHASE,
SUCCESS_PHASE,
STAGE,
ACTIVE_PHASE.as_ref(),
),
STAGE,
comment,
)
.await,
FinalizeOutcome::Applied
) {
return false;
}
info!(
ticket = %ticket.id,
"{STAGE}: the project commands and every participant passed (≥ {threshold}/10)",
threshold = VERIFICATION_THRESHOLD,
);
let _ = crate::jobs::terminalize_job(&crate::session::store().conn, job_id).await;
true
}
#[must_use]
fn verdict_passes(verdict: &crate::Verdict) -> bool {
verdict.score >= VERIFICATION_THRESHOLD
}
async fn git_available(ws: &Workspace) -> bool {
crate::git::commands::git_is_installed().await
&& crate::git::commands::is_git_repo(ws.as_path())
}
async fn skip_review_for_identical_content(ticket: &Ticket, ws: &Workspace) -> bool {
if !git_available(ws).await {
return false;
}
match compute_review_skip(ticket, ws.as_path()).await {
Ok(skip) => skip,
Err(e) => {
warn!(
ticket = %ticket.id,
error = %e,
"Git status check failed for skip-review — dispatching the reviewers",
);
false
}
}
}
fn should_skip_review(
reviewed_head: Option<&str>,
reviewed_tree: Option<&str>,
current_head: Option<&str>,
current_tree: Option<&str>,
porcelain: &str,
) -> bool {
let (Some(base_head), Some(base_tree)) = (reviewed_head, reviewed_tree) else {
return false;
};
let (Some(head), Some(tree)) = (current_head, current_tree) else {
return false;
};
head == base_head && tree == base_tree && !has_unstaged_changes(porcelain)
}
async fn compute_review_skip(ticket: &Ticket, repo_path: &Path) -> anyhow::Result<bool> {
let porcelain = run_git_status(repo_path).await?;
let head = run_git_head(repo_path).await.ok();
let tree = run_git_write_tree(repo_path).await.ok();
if (head.is_none() || tree.is_none()) && ticket.reviewed_head.is_some() {
warn!(
ticket = %ticket.id,
head = head.is_some(),
tree = tree.is_some(),
"Could not compute full content identity — dispatching the reviewers",
);
}
Ok(should_skip_review(
ticket.reviewed_head.as_deref(),
ticket.reviewed_tree.as_deref(),
head.as_deref(),
tree.as_deref(),
&porcelain,
))
}
async fn working_tree_churn(repo_path: &Path) -> anyhow::Result<i64> {
let snapshot = run_git_worktree_snapshot(repo_path).await?;
if snapshot.unborn_head {
anyhow::bail!("Repository has no commits — no churn baseline");
}
Ok(snapshot.stats.added + snapshot.stats.removed)
}
pub(crate) async fn compute_reviewer_count(ticket: &Ticket, repo_path: &Path) -> usize {
let tiny = crate::pipeline::verdict::DEFAULT_REVIEW_COUNT_TINY_CHURN;
let low = crate::pipeline::verdict::DEFAULT_REVIEW_COUNT_LOW_CHURN;
let high = crate::pipeline::verdict::DEFAULT_REVIEW_COUNT_HIGH_CHURN;
let count = match working_tree_churn(repo_path).await {
Ok(total) => {
let base = crate::pipeline::verdict::review_base_from_signals(total, tiny, low, high);
debug!(
ticket = %ticket.id,
total_churn = total,
reviewer_base = base,
"Reviewer count calibration: base {base} from total churn",
);
crate::pipeline::verdict::review_agent_count(base, ticket.priority)
}
Err(e) => {
warn!(
ticket = %ticket.id,
error = %e,
"Could not compute working-tree churn — reviewer base defaults to 3",
);
3
}
};
count.max(1)
}
async fn record_reviewed_base(
ws: &Workspace,
ticket_id: &str,
base: Option<ReviewBase>,
transitioned: bool,
) {
if !transitioned || !git_available(ws).await {
return;
}
if let Err(e) = run_git_add_all(ws.as_path()).await {
warn!(
ticket = %ticket_id,
error = %e,
"Failed to stage changes after review — reviewed base not recorded",
);
return;
}
let Some(base) = base else {
warn!(
ticket = %ticket_id,
"Could not read the reviewed content before the round — reviewed base not recorded",
);
return;
};
if let Err(e) = super::board()
.set_reviewed_base(ticket_id, Some(&base.head), Some(&base.tree))
.await
{
warn!(
ticket = %ticket_id,
error = %e,
"Failed to record reviewed base — later rounds will re-review",
);
} else {
debug!(ticket = %ticket_id, "Recorded reviewed base after verification");
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::git::commands::MAX_UNTRACKED_SIZE;
use crate::util::test::init_temp_repo;
#[tokio::test]
async fn working_tree_churn_ignores_huge_binary_file_count() {
let (_dir, repo_path) = init_temp_repo();
std::fs::write(repo_path.join("a.rs"), b"fn foo() {\n bar();\n}\n").unwrap();
let size = usize::try_from(MAX_UNTRACKED_SIZE).unwrap() + 1;
std::fs::write(repo_path.join("big.bin"), vec![b'a'; size]).unwrap();
let churn = working_tree_churn(&repo_path).await.unwrap();
assert_eq!(churn, 3);
}
}