use std::fmt::Write as _;
use std::sync::Arc;
use crate::jobs::RowStatus;
use crate::retry::RetryExhausted;
use crate::util::{panic_message, scrub_credentials};
use crate::{
AnalysisVerdict, BlockerVerificationVerdict, ChatMessage, ChatRequest, ChatRequestMeta, Role,
Verdict, Workspace,
};
use super::raw_response_dump_section;
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(),
}
}
}
pub(crate) struct JointFailure {
pub dump: String,
}
pub(crate) struct JointRound {
pub stage: &'static str,
pub verdicts: Vec<JointVerdict>,
pub failures: Vec<JointFailure>,
pub issues: Vec<Vec<String>>,
pub grades: Vec<Vec<Option<crate::IssueGrade>>>,
}
impl JointRound {
#[must_use]
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,
}),
}
}
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 {
let summary = if round.n_valid() > 0 || round.failures.is_empty() {
"\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 participant 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_role(name: &str) -> Option<Role> {
match name {
crate::pipeline::analysis::STAGE => Some(Role::Analyst),
crate::pipeline::verification::STAGE | "QA" => Some(Role::Qa),
"Review" => Some(Role::Reviewer),
_ => None,
}
}
#[derive(Clone)]
pub(crate) enum ParallelVerdict {
NoResponse(String),
ParseFailed(RetryExhausted),
Verdict(Verdict),
Analysis(AnalysisVerdict),
BlockerVerification(BlockerVerificationVerdict),
}
impl ParallelVerdict {
#[must_use]
pub(crate) fn is_technical_failure(&self) -> bool {
matches!(self, Self::NoResponse(_) | Self::ParseFailed(_))
}
}
#[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_grouping(
stage: &'static str,
results: &[ParallelVerdict],
role: Role,
single_participant_round: bool,
ws: &Workspace,
ticket_id: &str,
ticket_title: &str,
) -> (JointRound, crate::consensus::RepairOutcome) {
let round = build_joint_round(stage, results);
let has_no_issues = round.has_no_issues();
if has_no_issues || single_participant_round {
(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]) -> 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,
issues,
grades,
}
}