use std::sync::Arc;
use crate::pipeline::board::Ticket;
use crate::prompt::{load_prompt, load_prompt_sections, substitute};
use crate::{Role, Workspace};
use super::{
AgentSlot, BoardStore, ExtractionMode, FinalizeOutcome, JointRound, ParallelVerdict,
TicketPhase, TransitionCtx, Write, agent_slot_from_roster_row, build_agent_slots,
build_round_grouping, debug, deserialize_verdict_outcome, guard_job_phase, info,
insert_round_slots, issue_grade, pause_freezing, render_joint_comment, reset_phase_attempt,
run_parallel_agents, warn, with_comment_and_transition,
};
const DEFAULT_PARALLEL_AGENT_COUNT: usize = 3;
pub(crate) const STAGE: &str = "Analysis";
fn load_ticket_analysis_angles() -> Vec<String> {
load_prompt_sections("analyze/ticket_angles.md")
}
pub(crate) async fn run(ticket: Arc<Ticket>, ws: Workspace, job_id: String) {
if guard_job_phase(&ticket.id, TicketPhase::Analysis, &job_id).await {
return;
}
dispatch_backlog_analysts(ticket, ws, &job_id).await;
}
async fn ensure_analysis_slots(
ticket: &Ticket,
job_id: &str,
prompt: &str,
count: usize,
) -> Option<Vec<AgentSlot>> {
let slots = build_agent_slots(
&ticket.id,
Role::Analyst,
prompt,
&load_ticket_analysis_angles(),
0,
count,
);
match insert_round_slots(job_id, &[(crate::jobs::AgentKind::Analyst, &slots)]).await {
Ok(()) => Some(slots),
Err(e) => {
warn!(
ticket = %ticket.id,
job = %job_id,
error = %e,
"Failed to write the analysis roster — round not started",
);
None
}
}
}
async fn append_analysis_slots(
ticket: &Ticket,
job_id: &str,
prompt: &str,
count: usize,
) -> anyhow::Result<Vec<AgentSlot>> {
let roster = crate::jobs::list_agents_for_job(&crate::session::store().conn, job_id).await?;
let next_idx = roster
.iter()
.filter_map(|r| r.idx)
.max()
.map_or(0, |m| m + 1);
let slots = build_agent_slots(&ticket.id, Role::Analyst, prompt, &[], next_idx, count);
insert_round_slots(job_id, &[(crate::jobs::AgentKind::Analyst, &slots)]).await?;
Ok(slots)
}
pub(crate) struct EscalationEntry {
pub text: String,
}
pub(crate) struct ResolvedBlocker {
pub text: String,
pub kind: crate::BlockerKind,
pub severity: crate::BlockerSeverity,
pub impact: String,
pub reasoning: String,
}
fn display_kind(kind: crate::BlockerKind) -> &'static str {
match kind {
crate::BlockerKind::MainPathBlocker => "main-path blocker",
crate::BlockerKind::RiskEdgeCase => "risk/edge-case",
}
}
fn display_severity(severity: crate::BlockerSeverity) -> &'static str {
match severity {
crate::BlockerSeverity::Low => "low",
crate::BlockerSeverity::Medium => "medium",
crate::BlockerSeverity::High => "high",
crate::BlockerSeverity::Critical => "critical",
}
}
pub(crate) fn apply_blocker_verification(
entries: &[EscalationEntry],
verdicts: &[&crate::BlockerVerificationVerdict],
) -> Vec<ResolvedBlocker> {
entries
.iter()
.enumerate()
.map(|(i, entry)| {
let mut severity = crate::BlockerSeverity::Low;
let mut kind = crate::BlockerKind::RiskEdgeCase;
let mut impact_parts: Vec<&str> = Vec::new();
let mut reasoning_parts: Vec<&str> = Vec::new();
for round in verdicts {
for item in &round.verdicts {
if item.index != i {
continue;
}
if item.severity > severity {
severity = item.severity;
}
if item.kind == crate::BlockerKind::MainPathBlocker {
kind = item.kind;
}
if !item.impact.trim().is_empty() {
impact_parts.push(item.impact.trim());
}
if !item.reasoning.trim().is_empty() {
reasoning_parts.push(item.reasoning.trim());
}
}
}
ResolvedBlocker {
text: entry.text.clone(),
kind,
severity,
impact: impact_parts.join(" — "),
reasoning: reasoning_parts.join(" — "),
}
})
.collect()
}
const BLOCKER_LIST_MARKER: &str = "Blocker list:";
const LEGACY_BLOCKER_LIST_MARKER: &str = "matching index.";
fn format_blocker_list(entries: &[EscalationEntry]) -> String {
let mut out = format!("{BLOCKER_LIST_MARKER}\n");
for (i, e) in entries.iter().enumerate() {
let text = e.text.replace('\r', "").replace('\n', " ");
let _ = writeln!(out, "{i}. {text}");
}
out
}
fn format_blocker_verification_report(blockers: &[ResolvedBlocker]) -> String {
let mut out = String::from("\n\n### Blocker verification\n");
let _ = writeln!(
out,
"{} flagged blocker(s) verified with substance grades:",
blockers.len()
);
for (i, b) in blockers.iter().enumerate() {
let _ = writeln!(
out,
"\n{}. {} — {} (severity: {})",
i + 1,
b.text,
display_kind(b.kind),
display_severity(b.severity),
);
let _ = write!(out, "\n Impact: {}", b.impact);
let _ = write!(out, "\n Reasoning: {}", b.reasoning);
}
out
}
fn append_escalation_report(
joint_comment: &mut String,
entries: &[EscalationEntry],
escalation_results: &[ParallelVerdict],
) {
let verification: Vec<&crate::BlockerVerificationVerdict> = escalation_results
.iter()
.filter_map(|r| match r {
ParallelVerdict::BlockerVerification(v) => Some(v),
_ => None,
})
.collect();
let dispatched = escalation_results.len();
let succeeded = verification.len();
if succeeded == 0 {
if dispatched > 0 {
joint_comment.push_str(
"\n\n### Blocker verification\n\
The verification round could not produce a verdict — the base-round \
blockers remain unverified.",
);
}
return;
}
let resolved = apply_blocker_verification(entries, &verification);
joint_comment.push_str(&format_blocker_verification_report(&resolved));
if succeeded < dispatched {
let _ = write!(
joint_comment,
"\n\nNote: only {succeeded} of {dispatched} verifier(s) returned a verdict; the \
failed verifier(s) did not participate."
);
}
}
fn format_grade_summary(results: &[ParallelVerdict]) -> String {
let mut clean = 0usize;
let mut minor = 0usize;
let mut major = 0usize;
let mut blocker = 0usize;
let mut missing = 0usize;
for r in results {
match r {
ParallelVerdict::Analysis(v) => {
if v.issues_detected.is_empty() {
clean += 1;
} else {
match v.issues_detected.iter().map(|a| a.grade).max() {
Some(crate::IssueGrade::Blocker) => blocker += 1,
Some(crate::IssueGrade::Major) => major += 1,
_ => minor += 1,
}
}
}
ParallelVerdict::NoResponse(_) | ParallelVerdict::ParseFailed(_) => missing += 1,
_ => {}
}
}
let total = results.len();
let description = [
(blocker, "flagged a blocker"),
(major, "found major issues"),
(minor, "found minor issues"),
(clean, "found no issues"),
(missing, "provided no analysis"),
]
.iter()
.filter(|&&(count, _)| count > 0)
.map(|&(count, label)| {
if count == total {
format!("All {label}")
} else {
format!("{count} {label}")
}
})
.collect::<Vec<_>>()
.join(", ");
format!("{total} analysts reviewed this ticket. {description}.")
}
fn escalation_groups(
round: &JointRound,
outcome: &crate::consensus::RepairOutcome,
) -> Vec<EscalationEntry> {
let table = crate::consensus::ItemTable::new(&round.issues);
let first_blocker = |members: &[crate::consensus::GroupingMember]| {
for member in members {
for id in member.ids() {
if issue_grade(round, &table, id) == Some(crate::IssueGrade::Blocker) {
return Some(id);
}
}
}
None
};
match outcome {
crate::consensus::RepairOutcome::Repaired { output, .. } => {
let mut entries = Vec::new();
for group in &output.groups {
if let Some(id) = first_blocker(&group.members) {
let member_text = table
.resolve(id)
.map(|(_, t)| t.to_string())
.unwrap_or_default();
entries.push(EscalationEntry {
text: format!("{}: {}", group.heading, member_text),
});
}
}
for member in &output.ungrouped {
if first_blocker(std::slice::from_ref(member)).is_some()
&& let Some((_, text)) = table.resolve(member.id)
{
entries.push(EscalationEntry {
text: text.to_string(),
});
}
}
entries
}
crate::consensus::RepairOutcome::Fallback => {
let mut entries = Vec::new();
for id in 0..table.len() {
if issue_grade(round, &table, id) == Some(crate::IssueGrade::Blocker)
&& let Some((_, text)) = table.resolve(id)
{
entries.push(EscalationEntry {
text: text.to_string(),
});
}
}
entries
}
}
}
fn parse_escalation_entries(task: &str) -> Vec<EscalationEntry> {
let anchor = task
.find(BLOCKER_LIST_MARKER)
.map(|p| (p, BLOCKER_LIST_MARKER.len()))
.or_else(|| {
task.find(LEGACY_BLOCKER_LIST_MARKER)
.map(|p| (p, LEGACY_BLOCKER_LIST_MARKER.len()))
});
let Some((start, marker_len)) = anchor else {
return Vec::new();
};
let mut entries = Vec::new();
for line in task[start + marker_len..].lines() {
let Some(text) = split_numbered(line.trim()) else {
continue;
};
entries.push(EscalationEntry {
text: text.to_string(),
});
}
entries
}
fn split_numbered(line: &str) -> Option<&str> {
let after_digits = line.trim_start_matches(|c: char| c.is_ascii_digit());
if after_digits.len() == line.len() {
return None;
}
let after_dot = after_digits.strip_prefix('.')?;
let after_ws = after_dot.trim_start();
if after_ws.is_empty() {
return None;
}
Some(after_ws)
}
#[expect(clippy::too_many_arguments)]
#[must_use]
async fn maybe_escalate_analysis(
ticket: &Arc<Ticket>,
ws: &Workspace,
job_id: &str,
results: &mut Vec<ParallelVerdict>,
paused: &mut bool,
round: &JointRound,
outcome: &crate::consensus::RepairOutcome,
escalation_entries: &mut Vec<EscalationEntry>,
) -> bool {
let entries = escalation_groups(round, outcome);
if entries.is_empty() {
return true;
}
if crate::shutdown::aborting() {
return true;
}
if guard_job_phase(&ticket.id, TicketPhase::Analysis, job_id).await {
return false;
}
info!(
ticket = %ticket.id,
blockers = entries.len(),
"Base analysis flagged blocker groups — escalating with 2 additional analysts",
);
let escalation_task = substitute(
&load_prompt("analyze/blocker_verification.md"),
&[("{{blockers}}", &format_blocker_list(&entries))],
);
let extra_slots = match append_analysis_slots(ticket, job_id, &escalation_task, 2).await {
Ok(s) => s,
Err(e) => {
warn!(ticket = %ticket.id, error = %e, "Failed to append escalation slots — proceeding with base round");
return true;
}
};
if extra_slots.is_empty() {
return true;
}
let (extra, extra_paused) =
run_blocker_verification(ticket, ws, job_id, &entries, &extra_slots, false).await;
results.extend(extra);
*paused |= extra_paused;
*escalation_entries = entries;
true
}
async fn run_blocker_verification(
ticket: &Arc<Ticket>,
ws: &Workspace,
job_id: &str,
entries: &[EscalationEntry],
slots: &[AgentSlot],
resume: bool,
) -> (Vec<ParallelVerdict>, bool) {
let extraction_prompt = load_prompt("extraction/blocker_verification.md");
let blockers_arc =
Arc::<[String]>::from(entries.iter().map(|e| e.text.clone()).collect::<Vec<_>>());
run_parallel_agents(
ticket,
ws,
Role::Analyst,
&extraction_prompt,
ExtractionMode::BlockerVerification {
blockers: blockers_arc,
},
job_id,
slots,
TicketPhase::Analysis,
resume,
)
.await
}
async fn dispatch_backlog_analysts(ticket: Arc<Ticket>, ws: Workspace, job_id: &str) {
let message = super::analyst_task_prompt(&ticket);
let conn = &crate::session::store().conn;
let Some(roster) = super::read_roster_or_bail(&ticket.id, job_id).await else {
return;
};
if roster.is_empty() {
let Some(slots) =
ensure_analysis_slots(&ticket, job_id, &message, DEFAULT_PARALLEL_AGENT_COUNT).await
else {
return;
};
run_analysis_round(&ticket, &ws, job_id, &slots, false).await;
return;
}
let base_count = i64::try_from(DEFAULT_PARALLEL_AGENT_COUNT)
.expect("DEFAULT_PARALLEL_AGENT_COUNT fits in i64");
let escalation_rows: Vec<&crate::jobs::AgentRow> = roster
.iter()
.filter(|r| r.idx.unwrap_or(0) >= base_count)
.collect();
if escalation_rows.is_empty() {
let base_slots: Vec<AgentSlot> = roster.iter().map(agent_slot_from_roster_row).collect();
let not_done: Vec<String> = base_slots
.iter()
.filter(|s| s.status != crate::jobs::RowStatus::Done)
.map(|s| s.agent_id.clone())
.collect();
if let Err(e) = crate::jobs::rearm_roster_launched(conn, job_id, ¬_done).await {
warn!(ticket = %ticket.id, error = %e, "Failed to re-arm resumed analysis base slots");
}
run_analysis_round(&ticket, &ws, job_id, &base_slots, true).await;
} else {
let base_results: Vec<ParallelVerdict> = roster
.iter()
.filter(|r| r.idx.unwrap_or(0) < base_count)
.map(|r| deserialize_verdict_outcome(r.outcome.as_deref().unwrap_or("")))
.collect();
resume_escalation_round(&ticket, &ws, job_id, base_results).await;
}
}
async fn finalize_analysis_round_with_grouping(
ws: &Workspace,
ticket: &Ticket,
base_results: &[ParallelVerdict],
escalation_results: &[ParallelVerdict],
escalation_entries: &[EscalationEntry],
job_id: &str,
) {
let (round, outcome) = build_round_grouping(
STAGE,
base_results,
Role::Analyst,
false,
ws,
&ticket.id,
&ticket.title,
)
.await;
finalize_analysis_round(
ticket,
base_results,
escalation_results,
&round,
&outcome,
escalation_entries,
job_id,
)
.await;
}
async fn resume_escalation_round(
ticket: &Arc<Ticket>,
ws: &Workspace,
job_id: &str,
base_results: Vec<ParallelVerdict>,
) {
let conn = &crate::session::store().conn;
let roster = match crate::jobs::list_agents_for_job(conn, job_id).await {
Ok(roster) => roster,
Err(e) => {
warn!(ticket = %ticket.id, error = %e, "Failed to read escalation roster — retreating to base-only finalize");
finalize_analysis_round_with_grouping(ws, ticket, &base_results, &[], &[], job_id)
.await;
return;
}
};
let base_count = i64::try_from(DEFAULT_PARALLEL_AGENT_COUNT)
.expect("DEFAULT_PARALLEL_AGENT_COUNT fits in i64");
let escalation_slots: Vec<AgentSlot> = roster
.iter()
.filter(|r| r.idx.unwrap_or(0) >= base_count)
.map(agent_slot_from_roster_row)
.collect();
if escalation_slots.is_empty() {
finalize_analysis_round_with_grouping(ws, ticket, &base_results, &[], &[], job_id).await;
return;
}
let entries = parse_escalation_entries(&escalation_slots[0].task);
if entries.is_empty() {
finalize_analysis_round_with_grouping(ws, ticket, &base_results, &[], &[], job_id).await;
return;
}
if crate::shutdown::aborting()
|| guard_job_phase(&ticket.id, TicketPhase::Analysis, job_id).await
{
return;
}
let not_done: Vec<String> = escalation_slots
.iter()
.filter(|s| s.status != crate::jobs::RowStatus::Done)
.map(|s| s.agent_id.clone())
.collect();
if let Err(e) = crate::jobs::rearm_roster_launched(conn, job_id, ¬_done).await {
warn!(ticket = %ticket.id, error = %e, "Failed to re-arm resumed escalation roster slots");
}
let (extra, extra_paused) =
run_blocker_verification(ticket, ws, job_id, &entries, &escalation_slots, true).await;
if extra_paused {
pause_freezing(ticket, job_id).await;
return;
}
finalize_analysis_round_with_grouping(ws, ticket, &base_results, &extra, &entries, job_id)
.await;
}
async fn run_analysis_round(
ticket: &Arc<Ticket>,
ws: &Workspace,
job_id: &str,
slots: &[AgentSlot],
resume: bool,
) {
let extraction_prompt = load_prompt("extraction/analyst.md");
let (base_results, base_paused) = run_parallel_agents(
ticket,
ws,
Role::Analyst,
&extraction_prompt,
ExtractionMode::ScorelessVerdict,
job_id,
slots,
TicketPhase::Analysis,
resume,
)
.await;
if base_paused {
pause_freezing(ticket, job_id).await;
return;
}
let base_count = base_results.len();
let (round, outcome) = build_round_grouping(
STAGE,
&base_results,
Role::Analyst,
false,
ws,
&ticket.id,
&ticket.title,
)
.await;
let mut results = base_results;
let mut escalation_entries: Vec<EscalationEntry> = Vec::new();
let mut paused = false;
if !maybe_escalate_analysis(
ticket,
ws,
job_id,
&mut results,
&mut paused,
&round,
&outcome,
&mut escalation_entries,
)
.await
{
return;
}
if paused {
pause_freezing(ticket, job_id).await;
return;
}
let (base_results, escalation_results) = results.split_at(base_count);
finalize_analysis_round(
ticket,
base_results,
escalation_results,
&round,
&outcome,
&escalation_entries,
job_id,
)
.await;
}
async fn finalize_analysis_round(
ticket: &Ticket,
base_results: &[ParallelVerdict],
escalation_results: &[ParallelVerdict],
round: &JointRound,
outcome: &crate::consensus::RepairOutcome,
escalation_entries: &[EscalationEntry],
job_id: &str,
) {
if guard_job_phase(&ticket.id, TicketPhase::Analysis, job_id).await {
return;
}
if crate::shutdown::aborting() {
return;
}
process_analyst_verdicts(
ticket,
base_results,
escalation_results,
round,
outcome,
escalation_entries,
job_id,
)
.await;
if crate::shutdown::aborting() {
return;
}
if let Err(e) = crate::jobs::terminalize_job(&crate::session::store().conn, job_id).await {
warn!(job = %job_id, error = %e, "Failed to terminalize analysis job");
}
}
async fn reset_analysis_round(ticket: &Ticket, job_id: &str) {
info!(
ticket = %ticket.id,
"Backlog analysis produced no usable output — resetting for a fresh attempt",
);
let comment = "Backlog analysis produced no usable output.";
reset_phase_attempt(
ticket,
TicketPhase::Analysis,
job_id,
"analysis failure",
comment,
)
.await;
}
async fn process_analyst_verdicts(
ticket: &Ticket,
base_results: &[ParallelVerdict],
escalation_results: &[ParallelVerdict],
round: &JointRound,
outcome: &crate::consensus::RepairOutcome,
escalation_entries: &[EscalationEntry],
job_id: &str,
) {
let missing_analysis = base_results
.iter()
.filter(|r| r.is_technical_failure())
.count();
let dispatched = base_results.len();
let extracted_count = dispatched - missing_analysis;
if crate::shutdown::aborting() {
return;
}
if extracted_count == 0 {
reset_analysis_round(ticket, job_id).await;
return;
}
info!("{}", format_grade_summary(base_results));
let mut joint_comment = render_joint_comment(
round,
outcome,
&crate::consensus::ItemTable::new(&round.issues),
);
append_escalation_report(&mut joint_comment, escalation_entries, escalation_results);
if !matches!(
with_comment_and_transition(
TransitionCtx::notifying(
ticket,
TicketPhase::Analysis,
TicketPhase::Planning,
"Analyst",
Role::Analyst.as_str(),
),
async |tx| {
BoardStore::add_comment_tx(tx, &ticket.id, STAGE, &joint_comment).await?;
Ok(())
},
)
.await,
FinalizeOutcome::Applied
) {
return;
}
debug!(
ticket = %ticket.id,
nonempty_count = extracted_count,
"Backlog analysis complete — moved to planning ({extracted_count}/{dispatched} extracted)",
);
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn parse_escalation_entries_recover_marker_and_legacy_anchor() {
let current = "report the outcome for each blocker using the matching index.\n\n\
Blocker list:\n\
0. Score-less verdict model: the shared type is unspecified\n\
1. Crash-resume reconciliation: the divergence is inconsistent";
let entries = parse_escalation_entries(current);
assert_eq!(entries.len(), 2);
assert_eq!(
entries[0].text,
"Score-less verdict model: the shared type is unspecified"
);
assert_eq!(
entries[1].text,
"Crash-resume reconciliation: the divergence is inconsistent"
);
let legacy = "report the outcome for each blocker using the matching index.\n\n\
0. Score-less verdict model: the shared type is unspecified\n\
1. Crash-resume reconciliation: the divergence is inconsistent";
let entries = parse_escalation_entries(legacy);
assert_eq!(entries.len(), 2);
assert_eq!(
entries[0].text,
"Score-less verdict model: the shared type is unspecified"
);
assert!(parse_escalation_entries("no blocker list anywhere").is_empty());
}
}