use std::sync::Arc;
use crate::pipeline::board::Ticket;
use crate::prompt::{load_prompt, substitute};
use crate::{Agent, Role, Workspace};
use super::{
BoardStore, FinalizeOutcome, StageRunKind, TicketPhase, TransitionCtx, board,
bounce_to_development, clear_implementation_roster, determine_notify_policy, error,
guard_job_phase, guard_stage, info, list_new_or_untracked_files, notify_manager_system,
pause_freezing, pause_status_sentence, reset_phase_attempt, run_git_status, run_stage_agent,
sync_phase_job_task, warn, with_comment_and_transition,
};
pub(crate) async fn run(ticket: Arc<Ticket>, ws: Workspace, job_id: String) {
if guard_job_phase(&ticket.id, TicketPhase::InSanitation, &job_id).await {
return;
}
dispatch_sanitation(ticket, ws, &job_id).await;
}
async fn transition_ticket_to_done_no_comment(ticket: &Ticket, source: TicketPhase, reason: &str) {
let notify_policy = determine_notify_policy(&ticket.workspace_name, &ticket.id).await;
if matches!(
with_comment_and_transition(
TransitionCtx::new(
ticket,
source,
TicketPhase::Done,
notify_policy,
"Finalize",
Role::Sanitation.as_str(),
),
async |_tx| Ok(()),
)
.await,
FinalizeOutcome::Applied
) {
info!(ticket = %ticket.id, "{reason}");
let _ = crate::jobs::complete_ticket_phase_jobs(&crate::session::store().conn, &ticket.id)
.await;
}
}
async fn missing_repo_reason(ws: &Workspace) -> Option<&'static str> {
if !crate::git::commands::git_is_installed().await {
return Some("Git not installed — moving to Done without commit");
}
if !crate::git::commands::is_git_repo(ws.as_path()) {
return Some("Not a git repo — moving to Done without commit");
}
None
}
async fn finalize_sanitation_ticket(
ticket: &Ticket,
ws: &Workspace,
job_id: &str,
) -> Result<(), String> {
clear_implementation_roster(&crate::session::store().conn, job_id, &ticket.id).await;
let phase = TicketPhase::InSanitation;
if let Some(reason) = missing_repo_reason(ws).await {
transition_ticket_to_done_no_comment(ticket, phase, reason).await;
return Ok(());
}
let porcelain = run_git_status(ws.as_path())
.await
.map_err(|e| format!("{e:#}"))?;
if porcelain.trim().is_empty() {
transition_ticket_to_done_no_comment(
ticket,
phase,
"Clean working tree — moving to Done without commit",
)
.await;
return Ok(());
}
match crate::git::commands::run_git_commit(ws.as_path(), &ticket.title).await {
Ok(commit_info) => {
crate::git::commands::notify_git_commit(ws.as_path());
finalize_commit_and_transition(ticket, commit_info, phase).await;
Ok(())
}
Err(e) => Err(format!("{e:#}")),
}
}
async fn finish_or_stop(ticket: &Ticket, ws: &Workspace, job_id: &str) {
if let Err(cause) = finalize_sanitation_ticket(ticket, ws, job_id).await {
stop_sanitation_round(ticket, ws, job_id, &cause).await;
}
}
async fn stop_sanitation_round(ticket: &Ticket, ws: &Workspace, job_id: &str, cause: &str) {
error!(
ticket = %ticket.id,
error = %cause,
"Sanitation round could not finish — stopping the round and freezing the workspace"
);
let detail = crate::util::failure_detail(cause, "sanitation failure");
let comment = substitute(
&load_prompt("pipeline/sanitation_stop_comment.md"),
&[("{{failure_details}}", &detail)],
);
let was_frozen = matches!(
crate::workspace::store().get_by_name(&ws.name).await,
Ok(Some(live)) if live.paused
);
let pause_occurred = reset_phase_attempt(
ticket,
TicketPhase::InSanitation,
job_id,
"sanitation failure",
&comment,
)
.await;
let content = substitute(
&load_prompt("pipeline/sanitation_stop_notification.md"),
&[
("{{ticket_id}}", &ticket.id),
("{{failure_details}}", &detail),
(
"{{workspace_status}}",
&pause_status_sentence(pause_occurred || was_frozen),
),
],
);
notify_manager_system(&ws.name, content);
}
async fn finalize_commit_and_transition(
ticket: &Ticket,
commit_info: crate::git::commands::CommitInfo,
source: TicketPhase,
) {
let phase_label = source.as_ref();
crate::agent::registry::AGENT_REGISTRY.cancel_by_ticket_id(&ticket.id);
let notify_policy = determine_notify_policy(&ticket.workspace_name, &ticket.id).await;
let log_label = format!(
"finalize Done transition from {phase_label} ({})",
commit_info.short_hash(),
);
if matches!(
with_comment_and_transition(
TransitionCtx::new(
ticket,
source,
TicketPhase::Done,
notify_policy,
&log_label,
Role::Sanitation.as_str(),
),
async |tx| {
BoardStore::set_commit_info_tx(
tx,
&ticket.id,
&commit_info.hash,
commit_info.lines_added,
commit_info.lines_removed,
)
.await?;
Ok(())
},
)
.await,
FinalizeOutcome::Applied
) {
info!(ticket = %ticket.id, "Committed {}, moving to Done", commit_info.short_hash());
let _ = crate::jobs::complete_ticket_phase_jobs(&crate::session::store().conn, &ticket.id)
.await;
}
}
pub(crate) async fn finalize_sanitation_stage(
ticket: &Ticket,
agent: &Agent,
response: Option<&str>,
job_id: &str,
ws: &Workspace,
paused: bool,
) {
if guard_stage(
&ticket.id,
TicketPhase::InSanitation,
"Sanitation",
response,
job_id,
)
.await
{
return;
}
if paused {
pause_freezing(ticket, job_id).await;
return;
}
if response.is_none() {
warn!(
ticket = %ticket.id,
"Sanitation agent returned no output — resetting for a fresh attempt"
);
reset_phase_attempt(
ticket,
TicketPhase::InSanitation,
job_id,
"sanitation failure",
"Sanitation could not complete the round (the agent did not respond).",
)
.await;
return;
}
let extraction_prompt = crate::prompt::load_prompt("extraction/sanitation.md");
match agent
.extract_verdict::<crate::SanitationVerdict>(&extraction_prompt, None, None)
.await
{
Ok(verdict) => {
process_sanitation_verdict(ticket, job_id, verdict, ws).await;
}
Err(failure) => {
warn!(
ticket = %ticket.id,
error = %failure,
"Failed to extract sanitation verdict — resetting for a fresh attempt"
);
reset_phase_attempt(
ticket,
TicketPhase::InSanitation,
job_id,
"sanitation failure",
"Sanitation could not complete the round (the verdict could not be extracted).",
)
.await;
}
}
}
async fn dispatch_sanitation(ticket: Arc<Ticket>, ws: Workspace, job_id: &str) {
let untracked_files = match list_new_or_untracked_files(ws.as_path()).await {
Ok(files) if files.is_empty() => {
finish_or_stop(&ticket, &ws, job_id).await;
return;
}
Ok(files) => files.join("\n"),
Err(e) => {
let Some(reason) = missing_repo_reason(&ws).await else {
stop_sanitation_round(&ticket, &ws, job_id, &format!("{e:#}")).await;
return;
};
warn!(
ticket = %ticket.id,
error = %e,
reason,
"Failed to list untracked files — proceeding without a file list",
);
String::from("(could not list untracked files)")
}
};
let prompt = substitute(
&crate::prompt::load_prompt("sanitation.md"),
&[
("{{ticket_title}}", &ticket.title),
("{{ticket_description}}", &ticket.description),
("{{untracked_files}}", &untracked_files),
],
);
let conn = &crate::session::store().conn;
sync_phase_job_task(conn, job_id, &prompt).await;
run_stage_agent(&ticket, &ws, job_id, &prompt, StageRunKind::Sanitation).await;
}
async fn process_sanitation_verdict(
ticket: &Ticket,
job_id: &str,
verdict: crate::SanitationVerdict,
ws: &Workspace,
) {
if verdict.pass {
let passed_suffix = if verdict.garbage_files.is_empty() {
""
} else {
" (files reviewed)"
};
let comment = format!(
"🧹 Sanitation passed{passed_suffix}: {rationale}",
rationale = verdict.rationale
);
if let Err(e) = board()
.add_comment(&ticket.id, Role::Sanitation.as_str(), &comment)
.await
{
warn!(ticket = %ticket.id, error = %e, "Failed to record sanitation pass comment");
}
finish_or_stop(ticket, ws, job_id).await;
} else {
let garbage_list = verdict.garbage_files.join("\n- ");
let comment = substitute(
&load_prompt("pipeline/sanitation_failed_comment.md"),
&[
("{{garbage_list}}", &garbage_list),
("{{rationale}}", &verdict.rationale),
],
);
bounce_to_development(
ticket,
TicketPhase::InSanitation,
"Sanitation",
Role::Sanitation.as_str(),
Role::Sanitation.as_str(),
&comment,
job_id,
)
.await;
}
}