pub mod analysis;
pub mod board;
pub mod chronicle;
pub mod development;
pub mod diagnostics;
pub mod qa;
pub mod review;
pub mod sanitation;
#[cfg(test)]
mod tests;
pub(crate) mod verdict;
use std::collections::HashSet;
use std::fmt::Write;
use std::sync::{Arc, OnceLock};
use std::time::Duration;
use futures_util::future::join_all;
use tracing::{debug, error, info, warn};
use crate::agent::message_router;
use crate::agent::role::{DIAGNOSTICS_ROLE, SYSTEM_ROLE};
use crate::agent::{RETRY_EXHAUSTION_MARKER, run_agent};
use crate::db::TxGuard;
use crate::git::commands::{list_new_or_untracked_files, run_git_status};
use crate::pipeline::board::{BoardStore, Ticket, TicketPhase};
use crate::prompt::{load_prompt, load_prompt_sections, substitute};
use crate::session::manager_agent_id;
use crate::util::UnwrapPoison;
use crate::{Role, Workspace, WorkspaceStatus};
pub(crate) use verdict::stage_name;
pub(crate) use verdict::{
AgentSlot, ExtractionMode, JointRound, ParallelVerdict, build_round_grouping,
deserialize_verdict_outcome, issue_grade, process_verifier_verdicts, render_joint_comment,
round_member_failed, serialize_verdict_outcome, validate_blocker_verification,
validate_verdict_score,
};
pub(crate) use verdict::{QA_VI, REVIEWER_VI, VerifierInfo};
use development::finalize_engineer_stage;
use sanitation::finalize_sanitation_stage;
#[inline]
pub(crate) fn board() -> &'static BoardStore {
crate::pipeline::board::store()
}
static PHASE_BODIES_RUNNING: OnceLock<std::sync::Mutex<HashSet<String>>> = OnceLock::new();
fn phase_bodies_running() -> &'static std::sync::Mutex<HashSet<String>> {
PHASE_BODIES_RUNNING.get_or_init(|| std::sync::Mutex::new(HashSet::new()))
}
fn claim_phase_body(job_id: &str) -> bool {
phase_bodies_running()
.lock()
.unwrap_poison()
.insert(job_id.to_string())
}
fn release_phase_body(job_id: &str) {
phase_bodies_running().lock().unwrap_poison().remove(job_id);
}
struct PhaseBodyGuard(String);
impl Drop for PhaseBodyGuard {
fn drop(&mut self) {
release_phase_body(&self.0);
}
}
pub async fn run_management() {
let resumable = match crate::jobs::recover_from_restart().await {
Ok(r) => r,
Err(e) => {
error!(error = %e, "Boot recovery scan failed — proceeding without resume");
Vec::new()
}
};
if let Err(e) = crate::workspace::store()
.reclassify_analyzing_to_pending()
.await
{
warn!(error = %e, "Boot recovery: failed to reclassify stranded analyzing workspaces");
}
for stage in resumable {
let (job_id, workspace_name) = match &stage {
crate::jobs::ResumableJob::Research {
job_id,
workspace_name,
}
| crate::jobs::ResumableJob::Analyze {
job_id,
workspace_name,
}
| crate::jobs::ResumableJob::Implement {
job_id,
workspace_name,
}
| crate::jobs::ResumableJob::ResearchCleanup {
job_id,
workspace_name,
} => (job_id, workspace_name),
};
let Ok(Some(workspace)) = crate::users::resolve_workspace(workspace_name).await else {
warn!(
job = %job_id,
workspace = %workspace_name,
"Resume workspace unresolvable — job left in place (deleted only by explicit abandon)",
);
continue;
};
match stage {
crate::jobs::ResumableJob::Research { job_id, .. } => {
info!(job = %job_id, "Resuming research run at boot");
tokio::spawn(async move {
crate::tools::research::resume_research_run(&job_id, &workspace).await;
});
}
crate::jobs::ResumableJob::Analyze { job_id, .. } => {
info!(job = %job_id, "Resuming analyze round at boot");
tokio::spawn(async move {
crate::tools::analyze::resume_analyze_round(&job_id, &workspace).await;
});
}
crate::jobs::ResumableJob::Implement { job_id, .. } => {
info!(job = %job_id, "Resuming implement round at boot");
tokio::spawn(async move {
crate::tools::implement::resume_implement_round(&job_id, &workspace).await;
});
}
crate::jobs::ResumableJob::ResearchCleanup { job_id, .. } => {
info!(job = %job_id, "Resuming research cleanup agent at boot");
tokio::spawn(async move {
crate::research_cleanup::resume_research_cleanup(&job_id, &workspace).await;
});
}
}
}
let interval = Duration::from_millis(100);
loop {
if !crate::shutdown::sleep_or_shutdown_or_drain(interval).await {
break;
}
poll_round().await;
}
}
pub(crate) async fn poll_round() {
let workspaces = match crate::workspace::store().list().await {
Ok(ws_list) => ws_list,
Err(e) => {
error!(error = %e, "Failed to list workspaces");
return;
}
};
let tasks: Vec<_> = workspaces
.into_iter()
.map(|ws| {
tokio::spawn(async move {
process_single_workspace(ws).await;
})
})
.collect();
let results = join_all(tasks).await;
crate::util::log_join_failures(
results,
"Panic in workspace poll round — management loop continues",
"Workspace poll task was cancelled — management loop continues",
);
}
async fn process_single_workspace(ws: Workspace) {
let ws = match pickup_pending_workspace(&ws).await {
Some(claimed) => claimed,
None => ws,
};
if let Ok(Some(live)) = crate::workspace::store().get_by_name(&ws.name).await
&& live.paused
&& live.status == WorkspaceStatus::Ready
{
tracing::debug!(workspace = %ws.name, "Workspace paused mid-poll — skipping claim/dispatch for this round");
return;
}
run_claim_pipeline(&ws).await;
dispatch_working_phases(&ws).await;
}
async fn pickup_pending_workspace(ws: &Workspace) -> Option<Workspace> {
let (generation, discover_diagnostics) = pickup_claim(ws).await?;
info!(
workspace = %ws.name,
generation,
discover_diagnostics,
"Pickup: pending workspace claimed into discovery"
);
crate::workspace::spawn_workspace_discovery(ws, generation, discover_diagnostics);
crate::workspace::store()
.get_by_name(&ws.name)
.await
.ok()
.flatten()
}
async fn pickup_claim(ws: &Workspace) -> Option<(i64, bool)> {
if ws.status != WorkspaceStatus::Pending {
return None;
}
if !crate::config::provider_configured() {
return None;
}
if crate::workspace::pending_pickup_cooldown_active(&ws.name) {
return None;
}
let storage = crate::workspace::store();
let generation = match storage.claim_pending_for_discovery(&ws.name).await {
Ok(Some(generation)) => generation,
Ok(None) => return None,
Err(e) => {
warn!(
workspace = %ws.name,
error = %e,
"Pickup: failed to claim pending workspace — retrying next poll cycle"
);
return None;
}
};
let discover_diagnostics = match storage.get_by_name(&ws.name).await {
Ok(Some(fresh)) => fresh.diagnostics.is_none(),
Ok(None) | Err(_) => ws.diagnostics.is_none(),
};
Some((generation, discover_diagnostics))
}
async fn dispatch_working_phases(ws: &Workspace) {
let conn = &crate::session::store().conn;
let tickets = match board().list_working_tickets(&ws.name).await {
Ok(t) => t,
Err(e) => {
warn!(workspace = %ws.name, error = %e, "Failed to list working-phase tickets");
return;
}
};
for (phase, ticket) in tickets {
match crate::jobs::find_phase_job(conn, &ticket.id, phase).await {
Ok(Some(job)) => {
let has_agents = crate::jobs::job_has_launched_agents(conn, &job.id)
.await
.unwrap_or(true);
if has_agents {
continue;
}
if !claim_phase_body(&job.id) {
continue;
}
info!(ticket = %ticket.id, phase = %phase, job = %job.id, "Re-driving idle phase job");
spawn_phase_body(phase, Arc::new(ticket), ws.clone(), job.id);
continue;
}
Ok(None) => {}
Err(e) => {
warn!(ticket = %ticket.id, phase = %phase, error = %e, "Failed to check phase job");
continue;
}
}
let (task, role) = phase_task(&ticket, phase);
let job_id = crate::generate_id();
claim_phase_body(&job_id);
if let Err(e) = crate::jobs::spawn_job(
conn,
&job_id,
&task,
&ws.name,
"",
"",
role,
&[],
&crate::jobs::SpawnChild::Phase {
phase,
ticket_id: ticket.id.clone(),
},
None,
)
.await
{
release_phase_body(&job_id);
warn!(ticket = %ticket.id, phase = %phase, error = %e, "Failed to create phase job (concurrent creation is a no-op)");
continue;
}
info!(ticket = %ticket.id, phase = %phase, job = %job_id, "Created phase job — dispatching phase body");
spawn_phase_body(phase, Arc::new(ticket), ws.clone(), job_id);
}
}
fn phase_task(ticket: &Ticket, phase: TicketPhase) -> (String, Role) {
match phase {
TicketPhase::Analysis => (
crate::prompt::load_prompt(if ticket.reporter == Role::Maintainer.as_str() {
"analyze/maintainer_ticket.md"
} else {
"analyze/manager_ticket.md"
}),
Role::Analyst,
),
_ => (format!("Implement ticket {}", ticket.title), Role::Engineer),
}
}
async fn ticket_with_comments(ticket: Arc<Ticket>) -> Arc<Ticket> {
match board().get_ticket(&ticket.id).await {
Ok(Some(fresh)) => Arc::new(fresh),
Ok(None) | Err(_) => ticket,
}
}
fn spawn_phase_body(phase: TicketPhase, ticket: Arc<Ticket>, ws: Workspace, job_id: String) {
let log_job_id = job_id.clone();
let guard_job = job_id.clone();
tokio::spawn(async move {
let _guard = PhaseBodyGuard(guard_job);
let ticket = ticket_with_comments(ticket).await;
let panic_ticket = Arc::clone(&ticket);
let run: futures_util::future::BoxFuture<'static, ()> = match phase {
TicketPhase::Analysis => Box::pin(analysis::run(ticket, ws, job_id)),
TicketPhase::InDevelopment => Box::pin(development::run(ticket, ws, job_id)),
TicketPhase::InDiagnostics => Box::pin(diagnostics::run(ticket, ws, job_id)),
TicketPhase::InReview => Box::pin(review::run(ticket, ws, job_id)),
TicketPhase::InQa => Box::pin(qa::run(ticket, ws, job_id)),
TicketPhase::InSanitation => Box::pin(sanitation::run(ticket, ws, job_id)),
_ => {
error!(phase = %phase, "spawn_phase_body called for a non-working phase");
return;
}
};
if futures_util::FutureExt::catch_unwind(std::panic::AssertUnwindSafe(run))
.await
.is_ok()
{
return;
}
error!(phase = %phase, job = %log_job_id, "Phase body panicked — resetting for a fresh attempt");
let comment = format!("Pipeline phase {phase} crashed: the phase body panicked.");
reset_phase_attempt(
&panic_ticket,
phase,
&log_job_id,
"pipeline phase panic",
&comment,
)
.await;
});
}
async fn run_claim_pipeline(ws: &Workspace) {
let now = crate::db::now();
let grace_cutoff = (chrono::Utc::now() - BoardStore::BACKLOG_CLAIM_GRACE).to_rfc3339();
if board()
.backlog_claim_candidate_exists(&ws.name, &grace_cutoff)
.await
.unwrap_or(true)
&& let Ok(Some(ticket)) = board()
.claim_backlog_for_analysis(&ws.name, now.clone(), grace_cutoff)
.await
{
info!(ticket = %ticket.id, workspace = %ws.name, "Claimed Backlog → Analysis");
}
if board()
.queued_claim_candidate_exists(&ws.name)
.await
.unwrap_or(true)
&& let Ok(Some(ticket)) = board().claim_queued_for_development(&ws.name, now).await
{
info!(ticket = %ticket.id, workspace = %ws.name, "Claimed Queued → InDevelopment");
let _ = crate::jobs::upsert_session_pin(
&crate::session::store().conn,
&ticket.id,
&ticket.title,
crate::jobs::RowStatus::Launched,
Role::Engineer,
)
.await;
}
}
const PHASE_GATE_BAIL_REASON: &str = "ticket not in expected phase";
pub const MAX_BOUNCES: usize = 10;
pub const BOUNCE_BADGE_WARNING_THRESHOLD: usize = 5;
#[must_use]
async fn is_ticket_in_phase(ticket_id: &str, expected_phase: TicketPhase) -> bool {
match board().get_ticket_phase(ticket_id).await {
Ok(Some(phase)) => {
let ok = phase == expected_phase;
if !ok {
debug!(
ticket = %ticket_id,
expected_phase = %expected_phase,
actual = %phase,
"Ticket moved externally — bailing out",
);
}
ok
}
Ok(None) => {
debug!(ticket = %ticket_id, "Ticket not found — row missing (violates architecture invariant)");
false
}
Err(e) => {
warn!(ticket = %ticket_id, error = %e, "Failed to check ticket phase");
false
}
}
}
#[must_use]
async fn guard_job_phase(ticket_id: &str, expected: TicketPhase, job_id: &str) -> bool {
if !is_ticket_in_phase(ticket_id, expected).await {
if let Err(e) = crate::jobs::terminalize_job(&crate::session::store().conn, job_id).await {
warn!(job = %job_id, error = %e, "Failed to complete job after phase-moved");
}
return true;
}
false
}
#[derive(Debug, Copy, Clone, PartialEq, Eq)]
enum NotifyPolicy {
Notify,
Buffer,
}
#[derive(Debug)]
pub(crate) struct TransitionCtx<'t, 'l> {
ticket: &'t Ticket,
source: TicketPhase,
target: TicketPhase,
notify: NotifyPolicy,
log_label: &'l str,
breaker_trip: bool,
}
impl<'t, 'l> TransitionCtx<'t, 'l> {
fn new(
ticket: &'t Ticket,
source: TicketPhase,
target: TicketPhase,
notify: NotifyPolicy,
log_label: &'l str,
) -> Self {
Self {
ticket,
source,
target,
notify,
log_label,
breaker_trip: false,
}
}
fn notifying(
ticket: &'t Ticket,
source: TicketPhase,
target: TicketPhase,
log_label: &'l str,
) -> Self {
Self::new(ticket, source, target, NotifyPolicy::Notify, log_label)
}
pub(crate) fn buffered(
ticket: &'t Ticket,
source: TicketPhase,
target: TicketPhase,
log_label: &'l str,
) -> Self {
Self::new(ticket, source, target, NotifyPolicy::Buffer, log_label)
}
fn with_breaker(mut self, breaker_trip: bool) -> Self {
self.breaker_trip = breaker_trip;
self
}
}
#[derive(Debug, Copy, Clone, PartialEq, Eq)]
pub(crate) enum FinalizeOutcome {
Applied,
Moved,
Failed,
}
#[must_use]
async fn with_comment_and_transition<F>(
ctx: TransitionCtx<'_, '_>,
write_comments: F,
) -> FinalizeOutcome
where
F: AsyncFnOnce(&TxGuard<'_>) -> anyhow::Result<()>,
{
let outcome = match crate::db::with_tx_outcome(
&board().conn,
&ctx.ticket.id,
ctx.log_label,
async move |tx| {
write_comments(tx).await?;
BoardStore::transition_to_tx(tx, &ctx.ticket.id, Some(ctx.source), ctx.target).await
},
)
.await
{
Ok(true) => FinalizeOutcome::Applied,
Ok(false) => {
debug!(
ticket = %ctx.ticket.id,
"{}: ticket moved externally while in {} — finalization skipped (nothing written)",
ctx.log_label, ctx.source,
);
FinalizeOutcome::Moved
}
Err(e) => {
warn!(
ticket = %ctx.ticket.id,
error = %e,
"{}: transition to {} failed — ticket stuck in {}",
ctx.log_label, ctx.target, ctx.source,
);
FinalizeOutcome::Failed
}
};
if matches!(outcome, FinalizeOutcome::Applied) {
match ctx.notify {
NotifyPolicy::Notify => {
notify_ticket(ctx.ticket, ctx.source, ctx.target, ctx.breaker_trip).await;
}
NotifyPolicy::Buffer => {}
}
}
outcome
}
#[must_use]
pub(crate) async fn comment_and_transition(
ctx: TransitionCtx<'_, '_>,
role: &str,
text: &str,
) -> FinalizeOutcome {
let ticket = ctx.ticket;
with_comment_and_transition(ctx, async |tx| {
BoardStore::add_comment_tx(tx, &ticket.id, role, text).await?;
Ok(())
})
.await
}
async fn comment_and_transition_or_bail(
ctx: TransitionCtx<'_, '_>,
role: &str,
text: &str,
message: &str,
) {
let ticket_id = &ctx.ticket.id;
let target = ctx.target;
if !matches!(
comment_and_transition(ctx, role, text).await,
FinalizeOutcome::Applied
) {
return;
}
info!(ticket = %ticket_id, target = %target, "{message}");
}
#[must_use]
async fn resolve_ticket_workspace(ticket: &Ticket, log_label: &str) -> Option<Workspace> {
match crate::workspace::get_by_name(&ticket.workspace_name).await {
Ok(Some(ws)) => Some(ws),
Ok(None) => {
warn!(
ticket = %ticket.id,
workspace_name = %ticket.workspace_name,
"Workspace not found for ticket — {log_label}",
);
None
}
Err(e) => {
warn!(
ticket = %ticket.id,
workspace_name = %ticket.workspace_name,
error = %e,
"Failed to look up workspace for ticket — {log_label}",
);
None
}
}
}
async fn pause_freezing(ticket: &Ticket, job_id: &str) {
info!(
ticket = %ticket.id,
"Phase round paused — leaving job in place for the unpause re-drive"
);
if let Err(e) =
crate::jobs::interrupt_phase_job_roster(&crate::session::store().conn, job_id).await
{
warn!(ticket = %ticket.id, job = %job_id, error = %e, "Failed to mark paused agents on pause-freeze");
}
}
fn paused_workspace_sentence() -> &'static str {
"all in-flight work stops and no pipeline stage advances until the workspace is resumed"
}
pub(crate) async fn pause_workspace_on_failure(ticket: &Ticket, reason: &str) -> bool {
if crate::shutdown::aborting() {
return false;
}
let Some(ws) = resolve_ticket_workspace(ticket, "auto-pause skipped").await else {
return false;
};
if ws.paused {
return false;
}
match crate::workspace::store().set_paused(&ws.name, true).await {
Ok(()) => {
info!(
ticket = %ticket.id,
workspace = %ws.name,
reason,
"Workspace auto-paused after failure"
);
true
}
Err(e) => {
warn!(
ticket = %ticket.id,
workspace = %ws.name,
reason,
error = %e,
"Failed to auto-pause workspace after technical failure",
);
false
}
}
}
pub(crate) async fn reset_phase_attempt(
ticket: &Ticket,
phase: TicketPhase,
job_id: &str,
reason: &str,
comment: &str,
) {
crate::agent::registry::AGENT_REGISTRY.cancel_by_ticket_id(&ticket.id);
if phase.is_pipeline_occupied() {
pause_workspace_on_failure(ticket, reason).await;
}
if let Err(e) = board().add_comment(&ticket.id, SYSTEM_ROLE, comment).await {
warn!(ticket = %ticket.id, error = %e, "Failed to comment phase reset");
}
let _ = crate::jobs::terminalize_job(&crate::session::store().conn, job_id).await;
}
async fn last_comment_as_failure_details(ticket_id: &str) -> String {
match board().get_comments(ticket_id).await {
Ok(comments) => comments.last().map_or_else(
|| "(unknown failure reason)".to_string(),
|c| c.content.clone(),
),
Err(_) => "(unknown failure reason)".to_string(),
}
}
async fn notify_ticket(
ticket: &Ticket,
source: TicketPhase,
target_phase: TicketPhase,
breaker_trip: bool,
) {
let Some(ws) = resolve_ticket_workspace(ticket, "skipping notification").await else {
error!(
ticket = %ticket.id,
workspace_name = %ticket.workspace_name,
"Workspace resolution failed — notification skipped"
);
return;
};
let transition_log = format!(
"[{}] {}: {} → {}",
ticket.reporter,
ticket.id,
source.as_ref(),
target_phase.as_ref(),
);
if let Err(e) = crate::db::cdc::drain_once(&board().conn).await {
tracing::warn!(error = %e, "CDC flush before Manager notification failed");
}
let drained = crate::pipeline::chronicle::drain(&ws.name).await;
let mut message = substitute(
&load_prompt("pipeline/notification.md"),
&[
("{{ticket_id}}", &ticket.id),
("{{ticket_title}}", &ticket.title),
("{{ticket_phase}}", target_phase.as_ref()),
("{{transition_log}}", &transition_log),
("{{ticket_updates}}", &drained),
],
);
if target_phase == TicketPhase::Failed {
let failure_details = last_comment_as_failure_details(&ticket.id).await;
let workspace_status = if breaker_trip {
"Beware that all the other tickets have been moved back from Queued \
to Planning."
.to_string()
} else if ws.paused {
format!("The workspace is paused — {}.", paused_workspace_sentence())
} else {
"The workspace was not paused — remaining queued tickets may still be claimed."
.to_string()
};
let warning = substitute(
&load_prompt("pipeline/failure_notification.md"),
&[
("{{failure_details}}", &failure_details),
("{{workspace_status}}", &workspace_status),
],
);
message.push_str("\n\n");
message.push_str(&warning);
}
let agent_id = manager_agent_id(&ws.name);
message_router::route(
&agent_id,
message_router::AgentJob {
content: message,
workspace_name: ws.name,
user_name: String::new(),
channel: String::new(),
kind: message_router::MessageKind::TicketNotify,
role: crate::Role::Manager,
reply_target: None,
pending_job_id: None,
},
);
}
async fn register_running_agent(
job_id: &str,
agent_id: &str,
kind: crate::jobs::AgentKind,
warn_message: &str,
) -> tokio::sync::mpsc::UnboundedReceiver<message_router::AgentJob> {
let incoming_rx = message_router::register_agent(agent_id);
if let Err(e) = crate::jobs::upsert_job_agent(
&crate::session::store().conn,
job_id,
agent_id,
kind,
crate::jobs::RowStatus::Launched,
)
.await
{
warn!(
agent = agent_id,
job = job_id,
error = %e,
"{warn_message}",
);
}
incoming_rx
}
#[must_use]
fn agent_id(ticket_id: &str, idx: i64, suffix: &str, role: Role) -> String {
format!("ticket_{}_{}_{}_{}", ticket_id, idx, suffix, role.as_str())
}
#[must_use]
fn agent_slot_task(prompt: &str, angles: &[String], count: usize, global_idx: usize) -> String {
if angles.is_empty() {
prompt.to_string()
} else if count == 1 {
format!("{prompt}\n\n{}", angles.join("\n\n"))
} else {
format!("{prompt}\n\n{}", angles[global_idx % angles.len()])
}
}
#[must_use]
fn build_agent_slots(
ticket_id: &str,
role: Role,
prompt: &str,
angles: &[String],
start_idx: i64,
count: usize,
) -> Vec<AgentSlot> {
let suffix = crate::generate_suffix();
let mut slots = Vec::with_capacity(count);
let global_start = usize::try_from(start_idx).unwrap_or(usize::MAX);
for k in 0..count {
let idx = start_idx + i64::try_from(k).unwrap_or(i64::MAX);
let agent_id = agent_id(ticket_id, idx, &suffix, role);
let task = agent_slot_task(prompt, angles, count, global_start + k);
slots.push(AgentSlot {
idx,
agent_id,
task,
status: crate::jobs::RowStatus::Launched,
outcome: None,
});
}
slots
}
#[must_use]
fn agent_slot_from_roster_row(r: &crate::jobs::AgentRow) -> AgentSlot {
AgentSlot {
idx: r.idx.unwrap_or(0),
agent_id: r.agent_id.clone(),
task: r.task.clone(),
status: r
.status
.parse::<crate::jobs::RowStatus>()
.unwrap_or(crate::jobs::RowStatus::Failed),
outcome: r.outcome.clone(),
}
}
async fn insert_round_slots(job_id: &str, slots: &[AgentSlot], kind: crate::jobs::AgentKind) {
let conn = &crate::session::store().conn;
for slot in slots {
if let Err(e) = conn
.execute(
crate::jobs::AGENT_INSERT_SQL,
crate::jobs::agent_params(job_id, &slot.agent_id, kind, Some(slot.idx), &slot.task),
)
.await
{
warn!(
agent = %slot.agent_id,
job = %job_id,
error = %e,
"Failed to write round roster row",
);
}
}
}
async fn checkpoint_parallel_outcomes(
job_id: &str,
launched: &[&AgentSlot],
run_results: &[ParallelVerdict],
) -> std::collections::HashMap<String, ParallelVerdict> {
let conn = &crate::session::store().conn;
let mut by_agent = std::collections::HashMap::with_capacity(launched.len());
for (slot, result) in launched.iter().zip(run_results) {
let outcome = serialize_verdict_outcome(result);
let status = if matches!(
result,
ParallelVerdict::NoResponse(_) | ParallelVerdict::ParseFailed(_)
) {
crate::jobs::RowStatus::Failed
} else {
crate::jobs::RowStatus::Done
};
if let Err(e) =
crate::jobs::write_agent_outcome(conn, job_id, &slot.agent_id, status, Some(&outcome))
.await
{
warn!(
job = %job_id,
agent = %slot.agent_id,
error = %e,
"Failed to checkpoint agent outcome",
);
}
by_agent.insert(slot.agent_id.clone(), result.clone());
}
by_agent
}
fn assemble_parallel_results(
slots: &[AgentSlot],
by_agent: &std::collections::HashMap<String, ParallelVerdict>,
) -> Vec<ParallelVerdict> {
let mut results = Vec::with_capacity(slots.len());
for slot in slots {
if slot.status == crate::jobs::RowStatus::Done {
results.push(deserialize_verdict_outcome(
slot.outcome.as_deref().unwrap_or(""),
));
} else if let Some(r) = by_agent.get(slot.agent_id.as_str()) {
results.push(r.clone());
} else {
unreachable!("every non-Done slot is launched and recorded 1:1 in by_agent");
}
}
results
}
#[expect(clippy::too_many_arguments, clippy::too_many_lines)]
async fn run_parallel_agents(
ticket: &Arc<Ticket>,
ws: &Workspace,
role: Role,
extraction_prompt: &str,
extract_mode: ExtractionMode,
job_id: &str,
slots: &[AgentSlot],
expected_phase: TicketPhase,
resume: bool,
) -> (Vec<ParallelVerdict>, bool) {
let launched: Vec<&AgentSlot> = slots
.iter()
.filter(|s| s.status != crate::jobs::RowStatus::Done)
.collect();
let receivers: Vec<_> = launched
.iter()
.map(|s| message_router::register_agent(&s.agent_id))
.collect();
{
let members: Vec<_> = launched
.iter()
.zip(receivers)
.map(|(slot, rx)| {
let ticket = Arc::clone(ticket);
let ws = ws.clone();
let extraction_prompt = extraction_prompt.to_string();
let extract_mode = extract_mode.clone();
let agent_id = slot.agent_id.clone();
let task = slot.task.clone();
move |round: crate::agent::RoundOpts| async move {
if !is_ticket_in_phase(&ticket.id, expected_phase).await {
if let Some(notify) = &round.first_call_notify {
notify.notify_one();
}
message_router::unregister_agent(&agent_id);
return (
ParallelVerdict::NoResponse(PHASE_GATE_BAIL_REASON.to_string()),
false,
);
}
let has_session =
resume && crate::session::store().has_content(&agent_id).await;
let (agent, response) = run_agent(
agent_id.clone(),
role,
&ws,
Some(&ticket),
if has_session { "" } else { &task },
String::new(),
String::new(),
false,
Some(rx),
resume,
Some(round),
None,
None,
)
.await;
let response = response.unwrap_or_default();
if response.is_empty() {
if agent.is_paused_frozen() {
(
ParallelVerdict::NoResponse("agent paused".to_string()),
true,
)
} else {
let reason = agent.failure_reason("agent produced no response");
(
ParallelVerdict::NoResponse(crate::util::scrub_credentials(
&reason,
)),
false,
)
}
} else {
let verdict = match &extract_mode {
ExtractionMode::ScoreVerdict => agent
.extract_verdict::<crate::Verdict>(
&extraction_prompt,
Some(&validate_verdict_score),
None,
)
.await
.map(ParallelVerdict::Verdict),
ExtractionMode::ScorelessVerdict => agent
.extract_verdict::<crate::AnalysisVerdict>(
&extraction_prompt,
None,
None,
)
.await
.map(ParallelVerdict::Analysis),
ExtractionMode::BlockerVerification { blockers } => {
let blockers = std::sync::Arc::clone(blockers);
let validator = move |v: &crate::BlockerVerificationVerdict| {
validate_blocker_verification(v, &blockers)
};
agent
.extract_verdict::<crate::BlockerVerificationVerdict>(
&extraction_prompt,
Some(&validator),
None,
)
.await
.map(ParallelVerdict::BlockerVerification)
}
};
match verdict {
Ok(v) => (v, false),
Err(e) => (ParallelVerdict::ParseFailed(e), false),
}
}
}
})
.collect();
let handles = crate::agent::spawn_staggered_round(members, resume).await;
let mut run_results: Vec<ParallelVerdict> = Vec::with_capacity(handles.len());
let mut paused = false;
for handle in handles {
match handle.await {
Ok((v, p)) => {
paused |= p;
run_results.push(v);
}
Err(e) => run_results.push(round_member_failed(e)),
}
}
let by_agent = checkpoint_parallel_outcomes(job_id, &launched, &run_results).await;
(assemble_parallel_results(slots, &by_agent), paused)
}
}
pub(crate) fn raw_response_dump_section(failure: &crate::retry::RetryExhausted) -> String {
match failure.last_raw.as_deref() {
Some(text) if !text.trim().is_empty() => format!(
"Raw agent response (last attempt):\n```\n{}\n```",
crate::util::truncate_sandwich(
text,
crate::util::FAILURE_DETAIL_CAP,
"verdict response"
)
),
Some(_) => "Final attempt was a tool call — no text response produced.".to_string(),
None => format!(
"Extraction failed after {} attempt(s) — final failure: {} ({})",
failure.failures.len(),
failure.final_class.label(),
failure.detail,
),
}
}
fn load_verifier_angles(role: Role) -> Vec<String> {
match role {
Role::Reviewer => load_prompt_sections("review_angles.md"),
Role::Qa => load_prompt_sections("qa_angles.md"),
_ => Vec::new(),
}
}
#[derive(Clone, Copy)]
enum StageRunKind {
Engineer,
Sanitation,
}
impl StageRunKind {
fn role(self) -> Role {
match self {
Self::Engineer => Role::Engineer,
Self::Sanitation => Role::Sanitation,
}
}
fn agent_kind(self) -> crate::jobs::AgentKind {
match self {
Self::Engineer => crate::jobs::AgentKind::Engineer,
Self::Sanitation => crate::jobs::AgentKind::Sanitation,
}
}
}
async fn run_single_agent(
agent_id: String,
role: Role,
ws: &Workspace,
ticket: &Ticket,
message: &str,
incoming_rx: tokio::sync::mpsc::UnboundedReceiver<message_router::AgentJob>,
resume: bool,
) -> (crate::Agent, Option<String>) {
run_agent(
agent_id,
role,
ws,
Some(ticket),
message,
String::new(),
String::new(),
false,
Some(incoming_rx),
resume,
None,
None,
None,
)
.await
}
fn stage_drain_cut(ticket_id: &str, label: &str, response: Option<&str>) -> bool {
let drain_cut = response.is_none() && crate::shutdown::aborting();
if drain_cut {
info!(
ticket = %ticket_id,
"{label} round cut short by drain — job stays launched for boot resume",
);
}
drain_cut
}
async fn guard_stage(
ticket_id: &str,
phase: TicketPhase,
label: &str,
response: Option<&str>,
job_id: &str,
) -> bool {
if guard_job_phase(ticket_id, phase, job_id).await {
return true;
}
if stage_drain_cut(ticket_id, label, response) {
return true;
}
false
}
async fn run_stage_agent(
ticket: &Ticket,
ws: &Workspace,
job_id: &str,
task: &str,
kind: StageRunKind,
) {
let role = kind.role();
if let Err(e) = crate::jobs::upsert_session_pin(
&crate::session::store().conn,
&ticket.id,
task,
crate::jobs::RowStatus::Launched,
role,
)
.await
{
warn!(
ticket = %ticket.id,
job = %job_id,
error = %e,
"Failed to upsert {} session pin — session continuity across bounces/resets degraded",
role,
);
}
let agent_id = crate::jobs::session_pin_id(&ticket.id, role);
let incoming_rx = register_running_agent(
job_id,
&agent_id,
kind.agent_kind(),
match kind {
StageRunKind::Engineer => {
"Failed to register running engineer — stale agent already cancelled at dispatch, proceeding without roster registration"
}
StageRunKind::Sanitation => {
"Failed to register running sanitation agent — mid-run comments may not route"
}
},
)
.await;
let message = task.to_string();
let (agent, response) = run_single_agent(
agent_id,
kind.role(),
ws,
ticket,
&message,
incoming_rx,
false,
)
.await;
let paused = agent.is_paused_frozen();
match kind {
StageRunKind::Engineer => {
finalize_engineer_stage(ticket, &agent, response.as_deref(), job_id, ws, paused).await;
}
StageRunKind::Sanitation => {
finalize_sanitation_stage(ticket, &agent, response.as_deref(), job_id, ws, paused)
.await;
}
}
}
async fn sync_phase_job_task(conn: &crate::db::Connection, job_id: &str, task: &str) {
if let Err(e) = crate::jobs::update_phase_job_task(conn, job_id, task).await {
warn!(job = %job_id, error = %e, "Failed to sync phase job task");
}
}
async fn clear_implementation_roster(conn: &crate::db::Connection, job_id: &str, ticket_id: &str) {
if let Err(e) = crate::jobs::clear_launched_agents_for_job(conn, job_id).await {
warn!(ticket = %ticket_id, job = %job_id, error = %e, "Failed to clear running agents on stage handoff");
}
}
fn bounce_breaker_trip_comment() -> String {
let max = MAX_BOUNCES;
format!(
"Failed after {max} bounces — ticket bounced back too many times \
(circuit breaker, max: {max}). Ticket failed — Manager will triage."
)
}
#[must_use]
fn bounce_exhausted(bounce_count: i64) -> bool {
usize::try_from(bounce_count).unwrap_or(usize::MAX) >= MAX_BOUNCES
}
async fn drain_queued_siblings(ticket: &Ticket) {
match board()
.drain_queued_to_planning(&ticket.workspace_name)
.await
{
Ok(updated) if updated > 0 => {
info!(
tickets = updated,
workspace = %ticket.workspace_name,
"Moved {updated} Queued ticket(s) to Planning after bounce breaker trip",
);
}
Ok(_) => {
debug!(
workspace = %ticket.workspace_name,
"No Queued siblings to drain after bounce breaker trip",
);
}
Err(e) => {
warn!(
ticket = %ticket.id,
workspace = %ticket.workspace_name,
error = %e,
"Failed to move Queued tickets to Planning \
— breaker trip proceeds without moving siblings",
);
}
}
}
pub(crate) async fn bounce_to_development(
ticket: &Ticket,
source: TicketPhase,
log_label: &str,
drains_siblings: bool,
failure_role: &str,
failure_comment: &str,
job_id: &str,
) -> FinalizeOutcome {
let trip = bounce_exhausted(ticket.bounce_count);
let target = if trip {
TicketPhase::Failed
} else {
TicketPhase::InDevelopment
};
let notify = if trip {
NotifyPolicy::Notify
} else {
NotifyPolicy::Buffer
};
let trip_comment = trip.then(bounce_breaker_trip_comment);
let ctx = TransitionCtx::new(ticket, source, target, notify, log_label)
.with_breaker(trip && drains_siblings);
let outcome = with_comment_and_transition(ctx, async |tx| {
if let Some(comment) = &trip_comment {
BoardStore::add_comment_tx(tx, &ticket.id, SYSTEM_ROLE, comment).await?;
}
if !failure_comment.is_empty() {
BoardStore::add_comment_tx(tx, &ticket.id, failure_role, failure_comment).await?;
}
if !trip {
BoardStore::increment_bounce_count_tx(tx, &ticket.id).await?;
}
Ok(())
})
.await;
if matches!(outcome, FinalizeOutcome::Applied) {
let conn = &crate::session::store().conn;
if trip {
if drains_siblings {
drain_queued_siblings(ticket).await;
}
let _ = crate::jobs::terminalize_job(conn, job_id).await;
info!(
ticket = %ticket.id,
"Bounce circuit breaker tripped ({} bounces) — ticket failed",
MAX_BOUNCES,
);
} else {
let _ = crate::jobs::terminalize_job(conn, job_id).await;
info!(
ticket = %ticket.id,
target = %target,
"{log_label} failed — engineer re-dispatched via the puller on a fresh attempt",
);
}
}
outcome
}
async fn finalize_verifier_round(
ws: &Workspace,
ticket: &Ticket,
vi: VerifierInfo,
results: &[ParallelVerdict],
job_id: &str,
is_reviewer: bool,
) {
let transitioned = process_verifier_verdicts(ws, ticket, results, vi, job_id).await;
if !transitioned && crate::shutdown::aborting() {
info!(
ticket = %ticket.id,
"Verifier round cut short by drain — job stays launched for boot resume",
);
return;
}
if is_reviewer {
review::record_reviewed_base_after_review(
ws.as_path(),
&ticket.id,
review::git_available_for_review(ws, vi).await,
transitioned,
results,
)
.await;
}
}
async fn dispatch_verifiers(ticket: Arc<Ticket>, ws: Workspace, vi: VerifierInfo, job_id: String) {
if guard_job_phase(&ticket.id, vi.active_phase, &job_id).await {
return;
}
let is_reviewer = vi.role == Role::Reviewer;
if review::maybe_skip_review(&ticket, &ws, vi, &job_id).await {
return;
}
let extraction_prompt = crate::prompt::load_prompt(vi.extraction_prompt_path);
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, job = %job_id, error = %e, "Failed to read phase roster — bailing to preserve interrupted round");
return;
}
};
if roster.is_empty() {
fresh_dispatch_verifiers(ticket, ws, vi, job_id, is_reviewer).await;
return;
}
let slots: Vec<AgentSlot> = roster.iter().map(agent_slot_from_roster_row).collect();
let split = crate::jobs::split_slot_resume(&roster);
let count = slots.len();
let verifier_label = if count == 1 { "verifier" } else { "verifiers" };
info!(
ticket = %ticket.id,
role = %vi.role.as_str(),
count,
verifier_label,
"Resuming {count} parallel {verifier_label} from stored roster",
);
let not_done: Vec<String> = split.not_done.iter().map(|r| r.agent_id.clone()).collect();
if let Err(e) = crate::jobs::rearm_roster_launched(conn, &job_id, ¬_done).await {
warn!(ticket = %ticket.id, job = %job_id, error = %e, "Failed to re-arm resumed roster slots");
}
let (results, paused) = run_parallel_agents(
&ticket,
&ws,
vi.role,
&extraction_prompt,
ExtractionMode::ScoreVerdict,
&job_id,
&slots,
vi.active_phase,
true,
)
.await;
finish_verifier_dispatch(ticket, ws, vi, job_id, is_reviewer, results, paused).await;
}
async fn fresh_dispatch_verifiers(
ticket: Arc<Ticket>,
ws: Workspace,
vi: VerifierInfo,
job_id: String,
is_reviewer: bool,
) {
let engineer_response = ticket
.comments
.iter()
.rev()
.find(|c| c.role == Role::Engineer.as_str())
.map(|c| &c.content)
.map_or("(no output)", String::as_str);
let prompt = substitute(
&crate::prompt::load_prompt(vi.prompt_template),
&[("{{agent_response}}", engineer_response)],
);
let extraction_prompt = crate::prompt::load_prompt(vi.extraction_prompt_path);
let repo_path = ws.as_path();
let count = if is_reviewer {
review::compute_reviewer_count(&ticket, repo_path).await
} else {
qa::QA_PARALLEL_AGENT_COUNT
};
let verifier_label = if count == 1 { "verifier" } else { "verifiers" };
info!(
ticket = %ticket.id,
role = %vi.role.as_str(),
count,
verifier_label,
"Dispatching {count} parallel {verifier_label}",
);
let conn = &crate::session::store().conn;
sync_phase_job_task(conn, &job_id, &prompt).await;
let slots = build_agent_slots(
&ticket.id,
vi.role,
&prompt,
&load_verifier_angles(vi.role),
0,
count,
);
insert_round_slots(&job_id, &slots, crate::jobs::AgentKind::Verifier).await;
let (results, paused) = run_parallel_agents(
&ticket,
&ws,
vi.role,
&extraction_prompt,
ExtractionMode::ScoreVerdict,
&job_id,
&slots,
vi.active_phase,
false,
)
.await;
finish_verifier_dispatch(ticket, ws, vi, job_id, is_reviewer, results, paused).await;
}
async fn finish_verifier_dispatch(
ticket: Arc<Ticket>,
ws: Workspace,
vi: VerifierInfo,
job_id: String,
is_reviewer: bool,
results: Vec<ParallelVerdict>,
paused: bool,
) {
if guard_job_phase(&ticket.id, vi.active_phase, &job_id).await {
return;
}
if paused {
pause_freezing(&ticket, &job_id).await;
return;
}
finalize_verifier_round(&ws, &ticket, vi, &results, &job_id, is_reviewer).await;
}
async fn determine_notify_policy(workspace_name: &str, ticket_id: &str) -> NotifyPolicy {
match board()
.has_active_tickets_excluding(workspace_name, ticket_id)
.await
{
Ok(true) => {
debug!(
ticket = %ticket_id,
workspace = %workspace_name,
"Other active tickets remain — buffering Done notification",
);
NotifyPolicy::Buffer
}
Ok(false) => NotifyPolicy::Notify,
Err(e) => {
warn!(
ticket = %ticket_id,
workspace = %workspace_name,
error = %e,
"Failed to check active tickets — notifying to be safe",
);
NotifyPolicy::Notify
}
}
}