use std::fmt::Write as _;
use std::sync::Arc;
use crate::jobs::RowStatus;
use crate::pipeline::board::Ticket;
use crate::retry::RetryExhausted;
use crate::util::{panic_message, scrub_credentials};
use crate::{
AnalysisVerdict, BlockerVerificationVerdict, ChatMessage, ChatRequest, ChatRequestMeta, Role,
Verdict, Workspace,
};
use super::{
FinalizeOutcome, TicketPhase, TransitionCtx, bounce_to_development, comment_and_transition,
info, raw_response_dump_section, reset_phase_attempt,
};
pub(crate) const DEFAULT_REVIEW_COUNT_TINY_CHURN: i64 = 200;
pub(crate) const DEFAULT_REVIEW_COUNT_LOW_CHURN: i64 = 1000;
pub(crate) const DEFAULT_REVIEW_COUNT_HIGH_CHURN: i64 = 3000;
pub(crate) enum JointVerdict {
Score { verdict: crate::Verdict },
Graded { verdict: AnalysisVerdict },
}
impl JointVerdict {
#[must_use]
pub(crate) fn is_empty_issues(&self) -> bool {
match self {
Self::Score { verdict, .. } => verdict.issues_detected.is_empty(),
Self::Graded { verdict, .. } => verdict.issues_detected.is_empty(),
}
}
#[must_use]
pub(crate) fn passes(&self, threshold: u8) -> bool {
match self {
Self::Score { verdict, .. } => verdict.score >= threshold,
Self::Graded { .. } => true,
}
}
}
pub(crate) struct JointFailure {
pub dump: String,
}
pub(crate) struct JointRound {
pub stage: &'static str,
pub verdicts: Vec<JointVerdict>,
pub failures: Vec<JointFailure>,
pub threshold: u8,
pub issues: Vec<Vec<String>>,
pub grades: Vec<Vec<Option<crate::IssueGrade>>>,
}
impl JointRound {
#[must_use]
pub fn n_valid(&self) -> usize {
self.verdicts.len()
}
#[must_use]
pub fn has_no_issues(&self) -> bool {
self.verdicts.iter().all(JointVerdict::is_empty_issues)
}
}
const PIPELINE_GROUPING_MAX_TOKENS: u32 = 16_000;
fn synthesis_request(
round: &JointRound,
role: Role,
ws: &Workspace,
items: &[Vec<String>],
) -> ChatRequest {
let system = format!(
"{}\n\n{}",
crate::prompt::load_prompt("synthesis/synthesis.md"),
crate::prompt::load_prompt("synthesis/grouping_contradictions.md"),
);
let material = crate::consensus::numbered_items_material(items);
let user = format!(
"{}\n\nStage: {}\nAgent issues (id-numbered):\n{}",
crate::prompt::load_prompt("synthesis/synthesis_input.md"),
round.stage,
material,
);
let derived = crate::agent::role_chat_params(role);
ChatRequest {
messages: vec![ChatMessage::system(&system), ChatMessage::user(&user)],
tools: None,
model: derived.model,
max_tokens: Some(PIPELINE_GROUPING_MAX_TOKENS),
reasoning_effort: derived.reasoning_effort,
provider_order: derived.provider_order,
meta: Some(ChatRequestMeta {
purpose: "synthesis",
agent_id: format!("verdict_{}", crate::generate_suffix()),
role: role.as_str().to_string(),
workspace: ws.name.clone(),
ticket_id: None,
}),
}
}
pub(crate) async fn run_synthesis(
round: &JointRound,
role: Role,
ws: &Workspace,
ticket_id: &str,
ticket_title: &str,
) -> crate::consensus::RepairOutcome {
let request = synthesis_request(round, role, ws, &round.issues);
crate::consensus::run_grouping_repair(
ws,
"synthesis",
request,
&round.issues,
Some(crate::agent::registry::ParentKey::Ticket(
ticket_id.to_string(),
)),
Some(ticket_title.to_string()),
)
.await
}
#[must_use]
pub(crate) fn render_joint_comment(
round: &JointRound,
outcome: &crate::consensus::RepairOutcome,
table: &crate::consensus::ItemTable<'_>,
) -> String {
let mut out = String::new();
let has_issues = table.len() > 0;
if has_issues {
match outcome {
crate::consensus::RepairOutcome::Repaired { output, references } => {
for group in &output.groups {
let _ = write!(out, "\n\n**{}**", group.heading);
if group.contradiction {
out.push_str(" — DISPUTED");
}
for member in &group.members {
let _ = write!(out, "\n- {}", member_bullet(round, table, member));
}
}
out.push_str(&crate::consensus::render_ungrouped_section(
output,
references,
|member, disputed| member_bullet(round, table, member) + disputed,
));
}
crate::consensus::RepairOutcome::Fallback => {
out.push_str("\n\n**Issues**");
for id in 0..table.len() {
if let Some((_, text)) = table.resolve(id) {
let line = match issue_grade(round, table, id) {
Some(grade) => format!("{}: {}", grade.as_str(), text),
None => text.to_string(),
};
let _ = write!(out, "\n- {line}");
}
}
}
}
}
match outcome {
crate::consensus::RepairOutcome::Repaired { output, .. } => {
let summary = output.summary.trim();
if !summary.is_empty() {
out.push_str("\n\n### Summary");
let _ = write!(out, "\n{summary}");
}
}
crate::consensus::RepairOutcome::Fallback => {
if has_issues {
} else {
let clean = round.failures.is_empty()
&& round.verdicts.iter().all(|v| v.passes(round.threshold));
let summary = if clean || round.n_valid() > 0 {
"\n\n### Summary\nNo issues found.".to_string()
} else {
"\n\n### Summary\nNo issues to merge — no agent produced a verdict.".to_string()
};
out.push_str(&summary);
}
}
}
if !round.failures.is_empty() {
out.push_str("\n\n### Plain verifier responses");
for f in &round.failures {
let _ = write!(out, "\n- {}\n", f.dump);
}
}
crate::util::failure_detail(out.trim_start_matches('\n'), "joint verdict comment")
}
fn member_text(
table: &crate::consensus::ItemTable<'_>,
member: &crate::consensus::GroupingMember,
) -> String {
table.resolve(member.id).map_or_else(
|| format!("<unknown item id {}>", member.id),
|(_, text)| text.to_string(),
)
}
#[must_use]
pub(crate) fn issue_grade(
round: &JointRound,
table: &crate::consensus::ItemTable<'_>,
id: usize,
) -> Option<crate::IssueGrade> {
let (agent, item) = table.resolve_index(id)?;
round
.grades
.get(agent)
.and_then(|g| g.get(item))
.copied()
.flatten()
}
fn member_bullet(
round: &JointRound,
table: &crate::consensus::ItemTable<'_>,
member: &crate::consensus::GroupingMember,
) -> String {
match issue_grade(round, table, member.id) {
Some(grade) => format!("{}: {}", grade.as_str(), member_text(table, member)),
None => member_text(table, member),
}
}
#[must_use]
pub(crate) fn review_base_from_signals(
total_churn: i64,
tiny_churn: i64,
low_churn: i64,
high_churn: i64,
) -> usize {
if total_churn < tiny_churn {
1
} else if total_churn < low_churn {
2
} else if total_churn < high_churn {
3
} else {
4
}
}
#[must_use]
pub(crate) fn review_agent_count(base: usize, priority: i64) -> usize {
if priority == 0 { base.max(2) } else { base }
}
#[must_use]
pub(crate) fn stage_name(role: Role) -> &'static str {
match role {
Role::Analyst => "Analysis",
Role::Reviewer => "Review",
Role::Qa => "QA",
_ => unreachable!("stage_name called with a non-verdict role"),
}
}
#[must_use]
pub(crate) fn stage_role(name: &str) -> Option<Role> {
match name {
"Analysis" => Some(Role::Analyst),
"Review" => Some(Role::Reviewer),
"QA" => Some(Role::Qa),
_ => None,
}
}
#[derive(Clone)]
pub(crate) enum ParallelVerdict {
NoResponse(String),
ParseFailed(RetryExhausted),
Verdict(Verdict),
Analysis(AnalysisVerdict),
BlockerVerification(BlockerVerificationVerdict),
}
#[derive(Clone)]
pub(crate) enum ExtractionMode {
ScoreVerdict,
ScorelessVerdict,
BlockerVerification { blockers: Arc<[String]> },
}
#[derive(Clone)]
pub(crate) struct AgentSlot {
pub idx: i64,
pub agent_id: String,
pub task: String,
pub status: RowStatus,
pub outcome: Option<String>,
}
pub(crate) fn round_member_failed(e: tokio::task::JoinError) -> ParallelVerdict {
let reason = scrub_credentials(&panic_message(&*e.into_panic()));
tracing::warn!(%reason, "round member task failed");
ParallelVerdict::NoResponse(reason)
}
pub(crate) fn validate_verdict_score(v: &Verdict) -> Result<(), String> {
if v.score <= 10 {
Ok(())
} else {
Err(format!("verdict score {} out of range [0,10]", v.score))
}
}
pub(crate) fn validate_blocker_verification(
v: &BlockerVerificationVerdict,
blockers: &[String],
) -> Result<(), String> {
if v.verdicts.is_empty() {
return Err("blocker verification returned no verdicts".to_string());
}
if v.verdicts.len() != blockers.len() {
return Err(format!(
"blocker verification returned {} verdicts for {} blockers",
v.verdicts.len(),
blockers.len()
));
}
let mut seen = vec![false; blockers.len()];
for item in &v.verdicts {
if item.index >= blockers.len() {
return Err(format!("blocker index {} out of range", item.index));
}
if seen[item.index] {
return Err(format!("duplicate blocker index {}", item.index));
}
seen[item.index] = true;
if item.reasoning.trim().is_empty() {
return Err(format!("blocker {} missing reasoning", item.index));
}
if item.impact.trim().is_empty() {
return Err(format!("blocker {} missing impact", item.index));
}
}
Ok(())
}
#[must_use]
pub(crate) fn serialize_verdict_outcome(result: &ParallelVerdict) -> String {
match result {
ParallelVerdict::Verdict(v) => serde_json::json!({ "verdict": v }).to_string(),
ParallelVerdict::Analysis(v) => serde_json::json!({ "verdict": v }).to_string(),
ParallelVerdict::NoResponse(reason) => {
serde_json::json!({ "no_response": reason }).to_string()
}
ParallelVerdict::ParseFailed(f) => {
serde_json::json!({ "parse_failed": raw_response_dump_section(f) }).to_string()
}
ParallelVerdict::BlockerVerification(v) => {
serde_json::json!({ "blocker_verification": v }).to_string()
}
}
}
#[must_use]
pub(crate) fn deserialize_verdict_outcome(outcome: &str) -> ParallelVerdict {
let Ok(v) = serde_json::from_str::<serde_json::Value>(outcome) else {
return ParallelVerdict::NoResponse("unreadable stored outcome".to_string());
};
if let Some(verdict) = v.get("verdict") {
if verdict.get("score").is_some() {
if let Ok(vv) = serde_json::from_value::<crate::Verdict>(verdict.clone()) {
return ParallelVerdict::Verdict(vv);
}
} else if let Ok(av) = serde_json::from_value::<crate::AnalysisVerdict>(verdict.clone()) {
return ParallelVerdict::Analysis(av);
}
ParallelVerdict::NoResponse("unreadable stored verdict".to_string())
} else if let Some(r) = v.get("no_response").and_then(serde_json::Value::as_str) {
ParallelVerdict::NoResponse(r.to_string())
} else if let Some(p) = v.get("parse_failed").and_then(serde_json::Value::as_str) {
ParallelVerdict::NoResponse(p.to_string())
} else if let Some(v) = v.get("blocker_verification") {
match serde_json::from_value(v.clone()) {
Ok(bv) => ParallelVerdict::BlockerVerification(bv),
Err(_) => {
ParallelVerdict::NoResponse("unreadable stored blocker verification".to_string())
}
}
} else {
ParallelVerdict::NoResponse("unrecognized stored outcome".to_string())
}
}
pub(crate) async fn build_round_joint_comment(
stage: &'static str,
results: &[ParallelVerdict],
threshold: u8,
role: Role,
ws: &Workspace,
ticket_id: &str,
ticket_title: &str,
) -> String {
let (round, outcome) =
build_round_grouping(stage, results, threshold, role, ws, ticket_id, ticket_title).await;
render_joint_comment(
&round,
&outcome,
&crate::consensus::ItemTable::new(&round.issues),
)
}
pub(crate) async fn build_round_grouping(
stage: &'static str,
results: &[ParallelVerdict],
threshold: u8,
role: Role,
ws: &Workspace,
ticket_id: &str,
ticket_title: &str,
) -> (JointRound, crate::consensus::RepairOutcome) {
let round = build_joint_round(stage, results, threshold);
let has_no_issues = round.has_no_issues();
let single_verifier_verdict = matches!(role, Role::Reviewer | Role::Qa) && round.n_valid() == 1;
if has_no_issues || single_verifier_verdict {
(round, crate::consensus::RepairOutcome::Fallback)
} else {
let outcome = run_synthesis(&round, role, ws, ticket_id, ticket_title).await;
(round, outcome)
}
}
fn build_joint_round(
stage: &'static str,
results: &[ParallelVerdict],
threshold: u8,
) -> JointRound {
let mut verdicts: Vec<JointVerdict> = Vec::new();
let mut failures: Vec<JointFailure> = Vec::new();
let mut issues: Vec<Vec<String>> = vec![Vec::new(); results.len()];
let mut grades: Vec<Vec<Option<crate::IssueGrade>>> = vec![Vec::new(); results.len()];
for (i, r) in results.iter().enumerate() {
match r {
ParallelVerdict::Verdict(v) => {
verdicts.push(JointVerdict::Score { verdict: v.clone() });
issues[i].clone_from(&v.issues_detected);
grades[i] = vec![None; v.issues_detected.len()];
}
ParallelVerdict::Analysis(v) => {
verdicts.push(JointVerdict::Graded { verdict: v.clone() });
issues[i] = v.issues_detected.iter().map(|a| a.text.clone()).collect();
grades[i] = v.issues_detected.iter().map(|a| Some(a.grade)).collect();
}
ParallelVerdict::NoResponse(reason) => {
failures.push(JointFailure {
dump: reason.clone(),
});
}
ParallelVerdict::ParseFailed(f) => {
failures.push(JointFailure {
dump: scrub_credentials(&raw_response_dump_section(f)),
});
}
ParallelVerdict::BlockerVerification(_) => {}
}
}
JointRound {
stage,
verdicts,
failures,
threshold,
issues,
grades,
}
}
const REVIEW_QA_THRESHOLD: u8 = 9;
#[must_use]
fn verdict_passes(verdict: &crate::Verdict) -> bool {
verdict.score >= REVIEW_QA_THRESHOLD
}
#[derive(Copy, Clone)]
pub(crate) struct VerifierInfo {
pub(crate) role: Role,
pub(crate) log_label: &'static str,
pub(crate) success_phase: TicketPhase,
pub(crate) active_phase: TicketPhase,
pub(crate) prompt_template: &'static str,
pub(crate) extraction_prompt_path: &'static str,
}
pub(crate) const REVIEWER_VI: VerifierInfo = VerifierInfo {
role: Role::Reviewer,
log_label: "Reviewers",
success_phase: TicketPhase::InQa,
active_phase: TicketPhase::InReview,
prompt_template: "review.md",
extraction_prompt_path: "extraction/reviewer.md",
};
pub(crate) const QA_VI: VerifierInfo = VerifierInfo {
role: Role::Qa,
log_label: "QA",
success_phase: TicketPhase::InSanitation,
active_phase: TicketPhase::InQa,
prompt_template: "qa.md",
extraction_prompt_path: "extraction/qa.md",
};
pub(crate) async fn process_verifier_verdicts(
ws: &Workspace,
ticket: &Ticket,
results: &[ParallelVerdict],
verifier: VerifierInfo,
job_id: &str,
) -> bool {
let technical_failure = results.iter().any(|r| {
matches!(
r,
ParallelVerdict::NoResponse(_) | ParallelVerdict::ParseFailed(_)
)
});
let rework_failure = !technical_failure
&& results.iter().any(|r| match r {
ParallelVerdict::Verdict(v) => !verdict_passes(v),
_ => false,
});
if crate::shutdown::aborting() {
info!(
ticket = %ticket.id,
stage = %verifier.log_label,
"Verifier round cut short by drain — job stays launched for boot resume",
);
return false;
}
if technical_failure {
let comment = format!(
"{} could not complete the round (a verifier did not respond).",
verifier.log_label,
);
reset_phase_attempt(
ticket,
verifier.active_phase,
job_id,
verifier.log_label,
&comment,
)
.await;
return false;
}
let joint_comment = build_round_joint_comment(
stage_name(verifier.role),
results,
REVIEW_QA_THRESHOLD,
verifier.role,
ws,
&ticket.id,
&ticket.title,
)
.await;
if !rework_failure {
return apply_clean_verifier_round(ticket, verifier, &joint_comment, job_id).await;
}
let outcome = bounce_to_development(
ticket,
verifier.active_phase,
verifier.log_label,
true,
stage_name(verifier.role),
&joint_comment,
job_id,
)
.await;
matches!(outcome, FinalizeOutcome::Applied)
}
async fn apply_clean_verifier_round(
ticket: &Ticket,
verifier: VerifierInfo,
joint_comment: &str,
job_id: &str,
) -> bool {
if !matches!(
comment_and_transition(
TransitionCtx::buffered(
ticket,
verifier.active_phase,
verifier.success_phase,
verifier.log_label,
),
stage_name(verifier.role),
joint_comment,
)
.await,
FinalizeOutcome::Applied
) {
return false;
}
info!(
ticket = %ticket.id,
"{log_label}: all passed (≥ {threshold}/10)",
log_label = verifier.log_label,
threshold = REVIEW_QA_THRESHOLD,
);
let _ = crate::jobs::terminalize_job(&crate::session::store().conn, job_id).await;
true
}