use std::fmt::Write;
use std::sync::Arc;
use crate::pipeline::board::Ticket;
use crate::prompt::{load_prompt, substitute};
use crate::{Agent, Role, Workspace};
use super::{
RETRY_EXHAUSTION_MARKER, SYSTEM_ROLE, StageRunKind, TicketPhase, TransitionCtx, board,
comment_and_transition_or_bail, guard_job_phase, guard_stage, info, manager_agent_id,
message_router, pause_freezing, pause_status_sentence, pause_workspace_on_failure,
run_stage_agent, sync_phase_job_task, warn,
};
pub(crate) async fn run(ticket: Arc<Ticket>, ws: Workspace, job_id: String) {
if guard_job_phase(&ticket.id, TicketPhase::InDevelopment, &job_id).await {
return;
}
dispatch_engineer(ticket, ws, &job_id).await;
}
fn engineer_work_message(ticket: &Ticket) -> String {
let feedback: Vec<&str> = ticket
.comments
.iter()
.rposition(|c| c.role == Role::Engineer.as_str())
.map(|i| {
ticket
.comments
.iter()
.skip(i + 1)
.map(|c| c.content.as_str())
.collect()
})
.unwrap_or_default();
if feedback.is_empty() {
load_prompt("implement.md")
} else {
substitute(
&load_prompt("pipeline/bounce_feedback.md"),
&[("{{feedback}}", &feedback.join("\n---\n"))],
)
}
}
async fn dispatch_engineer(ticket: Arc<Ticket>, ws: Workspace, job_id: &str) {
let message = engineer_work_message(&ticket);
let conn = &crate::session::store().conn;
sync_phase_job_task(conn, job_id, &message).await;
run_stage_agent(&ticket, &ws, job_id, &message, StageRunKind::Engineer).await;
}
async fn engineer_comment_text(agent: &Agent, raw: &str) -> String {
let ticket_id = agent.ticket.as_ref().map_or("?", |t| t.id.as_str());
let policy = crate::retry::RetryPolicy::comment();
let extraction_prompt = load_prompt("extraction/engineer.md");
let summary = match agent
.extract_verdict::<crate::EngineerSummary>(&extraction_prompt, None, Some(&policy))
.await
{
Ok(summary) => summary,
Err(e) => {
warn!(
ticket = %ticket_id,
error = %e,
"Engineer summary extraction failed — using raw response for ticket comment"
);
return crate::util::scrub_credentials(raw);
}
};
let items: Vec<&str> = summary
.items
.iter()
.map(String::as_str)
.filter(|s| !s.trim().is_empty())
.collect();
if items.is_empty() {
warn!(
ticket = %ticket_id,
"Engineer summary extraction returned no usable items — using raw response for ticket comment"
);
return crate::util::scrub_credentials(raw);
}
let mut out = String::new();
for item in items {
let _ = write!(out, "\n- {}", item.replace('\n', " "));
}
let synopsis = summary.summary.as_deref().unwrap_or("").trim();
if !synopsis.is_empty() {
let _ = write!(out, "\n\n### Summary\n{synopsis}");
}
crate::util::failure_detail(&out, "engineer summary comment")
}
fn engineer_failure_comment(shutdown: bool, error: Option<&str>) -> String {
if shutdown {
return "Engineer failed: service shutting down — the run was interrupted \
by process shutdown."
.to_string();
}
let Some(detail) = error else {
return load_prompt("pipeline/engineer_failed.md");
};
let detail = crate::util::failure_detail(detail, "engineer failure");
if detail.contains(RETRY_EXHAUSTION_MARKER) {
format!("Engineer failed: LLM provider retry exhaustion.\n\n{detail}")
} else {
format!("Engineer failed.\n\n{detail}")
}
}
fn notify_engineer_pause(ws: &Workspace, failure_details: &str, paused: bool) {
let workspace_status = pause_status_sentence(paused);
let warning = substitute(
&load_prompt("pipeline/engineer_pause_notification.md"),
&[
("{{failure_details}}", failure_details),
("{{workspace_status}}", &workspace_status),
],
);
let agent_id = manager_agent_id(&ws.name);
message_router::route(
&agent_id,
message_router::AgentJob {
content: warning,
workspace_name: ws.name.clone(),
user_name: "system".to_string(),
channel: String::new(),
kind: message_router::MessageKind::UserMessage,
role: Role::Manager,
reply_target: None,
pending_job_id: None,
},
);
}
pub(crate) async fn finalize_engineer_stage(
ticket: &Ticket,
agent: &Agent,
response: Option<&str>,
job_id: &str,
ws: &Workspace,
paused: bool,
) {
if guard_stage(
&ticket.id,
TicketPhase::InDevelopment,
"Engineer",
response,
job_id,
)
.await
{
return;
}
if let Some(text) = response {
let comment_text = engineer_comment_text(agent, text).await;
comment_and_transition_or_bail(
TransitionCtx::buffered(
ticket,
TicketPhase::InDevelopment,
TicketPhase::InDiagnostics,
"Engineer",
Role::Engineer.as_str(),
),
Role::Engineer.as_str(),
&comment_text,
"Engineer finished — transitioned ticket",
)
.await;
let _ = crate::jobs::terminalize_job(&crate::session::store().conn, job_id).await;
return;
}
handle_engineer_failure(ticket, agent, job_id, ws, paused).await;
}
async fn handle_engineer_failure(
ticket: &Ticket,
agent: &Agent,
job_id: &str,
ws: &Workspace,
paused: bool,
) {
if paused {
pause_freezing(ticket, job_id).await;
return;
}
if agent.is_cancelled() && !crate::shutdown::aborting() {
info!(ticket = %ticket.id, "Engineer run interrupted by a code-driven cancellation — leaving the ticket for the replacement run");
return;
}
let pause_occurred = pause_workspace_on_failure(ticket, "engineer agent failure").await;
if crate::shutdown::aborting() {
info!(
ticket = %ticket.id,
"Engineer failure cut short by shutdown/drain after the pause — job stays launched for boot resume",
);
return;
}
let failure_comment = engineer_failure_comment(
crate::shutdown::shutdown_token().is_cancelled(),
agent.failure.as_deref(),
);
let conn = &crate::session::store().conn;
if let Err(e) = board()
.add_comment(&ticket.id, SYSTEM_ROLE, &failure_comment)
.await
{
warn!(ticket = %ticket.id, error = %e, "Failed to comment engineer hard failure");
}
let workspace_paused = ws.paused || pause_occurred;
notify_engineer_pause(ws, &failure_comment, workspace_paused);
info!(
ticket = %ticket.id,
"Engineer hard failure — workspace paused, ticket reset for a fresh development attempt"
);
let _ = crate::jobs::terminalize_job(conn, job_id).await;
}