use std::sync::Arc;
use super::heal_gate::{Unreachable, Verdict};
use super::heal_intake::{Checkout, HealTarget};
use super::heal_live::{parse_verdict, CoderRunner};
use super::heal_select::Candidate;
use super::heal_tick::{DeliverRefusal, RunFailure};
use super::merge::PrDeliveryOutcome;
use super::native_loop::TurnGenerator;
use super::provenance::SessionSeed;
use super::router::EngineChoice;
use super::rpc::{confirm_session, start_session, CoderSessionEntry, StartArgs};
use super::session::{CoderEventKind, CoderState};
use crate::session::ServerState;
const MAX_REVIEW_DIFF_BYTES: usize = 200_000;
pub const DEFAULT_ITEM_WALL_SECS: u64 = 45 * 60;
pub const PANEL_WALL_SECS: u64 = 10 * 60;
#[async_trait::async_trait]
pub trait Reviewer: Send + Sync {
fn model(&self) -> &str;
async fn review(&self, criteria: &str, diff: &str) -> Result<String, String>;
}
async fn session_entry(
state: &Arc<ServerState>,
id: &str,
) -> Result<Arc<CoderSessionEntry>, String> {
state
.coder_sessions
.lock()
.await
.get(id)
.cloned()
.ok_or_else(|| format!("coder session `{id}` is not registered"))
}
pub struct LiveCoderRunner {
pub state: Arc<ServerState>,
pub generator: Arc<dyn TurnGenerator>,
pub state_dir: std::path::PathBuf,
pub reviewers: Vec<Arc<dyn Reviewer>>,
pub max_wall_secs: u64,
pub max_iterations: Option<u32>,
pub engine: EngineChoice,
pub model: Option<String>,
pub routing_exclusions: Vec<String>,
pub canonical_model: Arc<dyn Fn(&str) -> String + Send + Sync>,
pub github: Arc<dyn super::merge::GitHubApi>,
}
pub fn checkout_path(target: &HealTarget) -> Result<std::path::PathBuf, String> {
match target.checkout.as_ref() {
Some(Checkout::Local(p)) => Ok(p.clone()),
Some(Checkout::Project(slug)) => Ok(super::project::load_project(slug)?.repo_path.clone()),
None => Err("target has no checkout; it is watch-only".into()),
}
}
fn configuration_run_failure(
session_id: &str,
detail: String,
failure_kind: Option<&str>,
) -> Option<RunFailure> {
(failure_kind == Some("configuration"))
.then(|| RunFailure::configuration_with_session(session_id, detail))
}
impl LiveCoderRunner {
async fn terminate(&self, session_id: &str, to: CoderState) {
let Ok(entry) = session_entry(&self.state, session_id).await else {
return;
};
let mut session = entry.session.lock().await;
if session.state.is_terminal() {
return;
}
if let Err(e) = session.transition(to, &entry.sink) {
tracing::warn!(session = %session_id, "self-heal could not close the session: {e}");
}
}
async fn poll_panel(&self, criteria: &str, diff: &str) -> (Vec<Verdict>, Vec<Unreachable>) {
let futures: Vec<_> = self
.reviewers
.iter()
.map(|r| async move {
let model = r.model().to_string();
match r.review(criteria, diff).await {
Ok(answer) => parse_verdict(&model, &answer),
Err(e) => Err(Unreachable { model, error: e }),
}
})
.collect();
let mut verdicts = Vec::new();
let mut unreachable = Vec::new();
for result in futures::future::join_all(futures).await {
match result {
Ok(v) => verdicts.push(v),
Err(u) => unreachable.push(u),
}
}
(verdicts, unreachable)
}
}
#[async_trait::async_trait]
impl CoderRunner for LiveCoderRunner {
async fn run(
&self,
target: &HealTarget,
item: &Candidate,
seed: &SessionSeed,
) -> Result<super::heal_tick::Attempt, RunFailure> {
let repo = checkout_path(target).map_err(RunFailure::early)?;
let start = start_session(
&self.state,
StartArgs {
distributed: false,
browser: false,
workers: Vec::new(),
repo,
intent: seed.as_str().to_string(),
engine: self.engine.clone(),
max_iterations: self.max_iterations,
state_dir: self.state_dir.clone(),
project: None,
model: self.model.clone(),
routing_exclusions: self.routing_exclusions.clone(),
repair_invokes: None,
transient_retries: None,
discussion_id: None,
base: None,
},
self.generator.clone(),
)
.await
.map_err(RunFailure::early)?;
let session_id = start["session_id"]
.as_str()
.ok_or_else(|| RunFailure::early("coder.start returned no session_id"))?
.to_string();
confirm_session(&self.state, &session_id, None)
.await
.map_err(|e| RunFailure::with_session(&session_id, e))?;
let entry = session_entry(&self.state, &session_id)
.await
.map_err(|e| RunFailure::with_session(&session_id, e))?;
{
let session = entry.session.lock().await;
let unrunnable = super::contract::baseline_cannot_run(&session.baseline);
if !unrunnable.is_empty() {
return Err(RunFailure::with_session(
&session_id,
format!(
"the derived contract cannot be evaluated for {}#{}: check(s) {} \
could not be run at all (the command does not exist here), so no \
change could ever make this contract green",
item.repo,
item.number,
unrunnable.join(", ")
),
));
}
if super::contract::baseline_gates_nothing(&session.baseline) {
return Err(RunFailure::with_session(
&session_id,
format!(
"every check of the contract derived for {}#{} already passes on \
unmodified code ({} check(s)), so it gates nothing and no change \
could turn it red-to-green. Either the issue's premise no longer \
holds — the behaviour it describes may already be fixed — or the \
derived checks do not exercise it. A human should decide which.",
item.repo,
item.number,
session.baseline.len()
),
));
}
}
let handle = entry.task.lock().unwrap().take();
if let Some(mut handle) = handle {
match tokio::time::timeout(
std::time::Duration::from_secs(self.max_wall_secs),
&mut handle,
)
.await
{
Ok(Ok(())) => {}
Ok(Err(e)) => {
return Err(RunFailure::with_session(
&session_id,
format!("coder session panicked: {e}"),
))
}
Err(_) => {
handle.abort();
self.terminate(&session_id, CoderState::Failed).await;
return Err(RunFailure::with_session(
&session_id,
format!(
"coder session exceeded its {}s ceiling for {}#{}",
self.max_wall_secs, item.repo, item.number
),
));
}
}
}
let (state, contract_detail, authors, failure_kind, ran_as, nominated) = {
let s = entry.session.lock().await;
let nominated = s.no_change_finding.as_ref().and_then(|finding| {
let green = !s.last_check_results.is_empty()
&& s.last_check_results.iter().all(|r| r.passed);
(!green).then(|| {
format!(
"the coder concluded no change is needed ({}: {}) without the \
outcome contract passing; a no-change conclusion needs a human",
finding.kind.as_str(),
finding.summary
)
})
});
(
s.state,
s.error
.clone()
.unwrap_or_else(|| "outcome contract evaluated".to_string()),
s.authored_by.clone(),
s.failure_kind.clone(),
s.engine.clone(),
nominated,
)
};
let red_nomination = nominated.is_some();
let contract_detail = nominated.unwrap_or(contract_detail);
let contract_passed = state == CoderState::NeedsApproval && !red_nomination;
if !contract_passed {
if let Some(failure) = configuration_run_failure(
&session_id,
contract_detail.clone(),
failure_kind.as_deref(),
) {
return Err(failure);
}
return Ok(super::heal_tick::Attempt {
contract_passed: false,
contract_detail,
panel_size: self.reviewers.len(),
verdicts: Vec::new(),
unreachable: Vec::new(),
session_id: session_id.clone(),
});
}
{
if authors.is_empty() && ran_as == EngineChoice::Native {
return Err(RunFailure::with_session(
&session_id,
format!(
"the native coder produced a change for {}#{} with no model \
recorded against it — the journal should name every model that \
completed a turn, so this session cannot be shown to be \
independent of the review panel",
item.repo, item.number
),
));
}
let seats: Vec<String> = self
.reviewers
.iter()
.map(|r| r.model().to_string())
.collect();
for author in &authors {
let Some(seat) = super::heal_review::coder_on_panel(author, &seats, |m| {
(self.canonical_model)(m)
}) else {
continue;
};
return Err(RunFailure::with_session(
&session_id,
format!(
"{} wrote part of this change for {}#{} and also sits on the review \
panel ({}) — a model cannot review its own output, and counting it \
would report an independence this gate does not have. The router \
chose it — nothing pinned a coder, or the pin is not what \
served this turn — so pin a coder that is not a seat (`heal.toml`'s \
`coder_model`, or `coder.toml`'s `model`), or drop that seat from \
`review_models`.",
author, item.repo, item.number, seat
),
));
}
}
let worktree = {
let s = entry.session.lock().await;
s.workspace_path.clone().ok_or_else(|| {
RunFailure::with_session(&session_id, "session reached approval with no worktree")
})?
};
let staged =
super::merge::stage_and_diff(&worktree, MAX_REVIEW_DIFF_BYTES).map_err(|e| {
RunFailure::with_session(
&session_id,
format!("could not read the change for review: {e}"),
)
})?;
if staged.truncated {
return Err(RunFailure::with_session(
&session_id,
format!(
"the change for {}#{} is {} bytes, larger than the {} a reviewer is \
given — a panel judging the tail of a patch is not a review",
item.repo, item.number, staged.full_bytes, MAX_REVIEW_DIFF_BYTES
),
));
}
if staged.changed_paths.is_empty() {
return Err(RunFailure::with_session(
&session_id,
format!(
"the outcome contract passed but the worktree is unchanged for {}#{} — there is nothing to review and nothing to deliver",
item.repo, item.number
),
));
}
let criteria = super::heal_gate::review_criteria(seed.as_str(), &contract_detail);
let (verdicts, unreachable) = match tokio::time::timeout(
std::time::Duration::from_secs(PANEL_WALL_SECS),
self.poll_panel(&criteria, &staged.patch),
)
.await
{
Ok(v) => v,
Err(_) => (
Vec::new(),
self.reviewers
.iter()
.map(|r| Unreachable {
model: r.model().to_string(),
error: format!("the panel did not answer within {PANEL_WALL_SECS}s"),
})
.collect(),
),
};
Ok(super::heal_tick::Attempt {
contract_passed: true,
contract_detail,
panel_size: self.reviewers.len(),
verdicts,
unreachable,
session_id,
})
}
async fn deliver(
&self,
target: &HealTarget,
item: &Candidate,
session_id: &str,
body: &str,
) -> Result<PrDeliveryOutcome, DeliverRefusal> {
let repo = checkout_path(target).map_err(DeliverRefusal::permanent)?;
self.github.auth_status().map_err(|e| {
DeliverRefusal::permanent(format!("gh is not logged in: {}", e.message))
})?;
let entry = session_entry(&self.state, session_id)
.await
.map_err(DeliverRefusal::permanent)?;
let (worktree, contract, intent) = {
let session = entry.session.lock().await;
(
session.workspace_path.clone().ok_or_else(|| {
DeliverRefusal::permanent("session has no worktree to deliver")
})?,
session
.contract
.clone()
.ok_or_else(|| DeliverRefusal::permanent("session has no outcome contract"))?,
session.intent.clone(),
)
};
let branch = delivery_branch(item);
let github = self.github.clone();
let base = target.base.clone();
let body = body.to_string();
let outcome = tokio::task::spawn_blocking(move || {
super::merge::deliver_pr_with(
super::merge::PrDelivery {
repo: &repo,
worktree: &worktree,
target_branch: &branch,
base_branch: &base,
draft: false,
intent: &intent,
contract: &contract,
body: &body,
provenance: None,
},
github.as_ref(),
)
})
.await
.map_err(|e| DeliverRefusal::retriable(format!("delivery task failed: {e}")))?
.map_err(refusal_for)?;
let mut session = entry.session.lock().await;
session.result_branch = Some(outcome.branch.clone());
entry.sink.emit(CoderEventKind::MergeCompleted {
branch: outcome.branch.clone(),
});
session
.transition(CoderState::Merged, &entry.sink)
.map_err(DeliverRefusal::permanent)?;
Ok(outcome)
}
async fn abandon(&self, session_id: &str) {
self.terminate(session_id, CoderState::Abandoned).await;
}
}
pub fn delivery_branch(item: &Candidate) -> String {
format!("car/heal/{}-{}", item.repo.replace('/', "-"), item.number)
}
fn refusal_for(failure: super::merge::DeliveryFailure) -> DeliverRefusal {
use super::merge::DeliveryFailure as F;
match failure {
F::Preflight { reason } => DeliverRefusal::permanent(reason),
F::Commit { reason } => DeliverRefusal::permanent(reason),
F::Push { reason, retriable } | F::Pr { reason, retriable } => {
if retriable {
DeliverRefusal::retriable(reason)
} else {
DeliverRefusal::permanent(reason)
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
struct Fixed {
model: &'static str,
answer: Result<String, String>,
}
#[async_trait::async_trait]
impl Reviewer for Fixed {
fn model(&self) -> &str {
self.model
}
async fn review(&self, _c: &str, _d: &str) -> Result<String, String> {
self.answer.clone()
}
}
#[test]
fn configuration_session_failure_is_typed_for_the_tick() {
let failure = configuration_run_failure(
"coder-1",
"no independent coder".into(),
Some("configuration"),
)
.expect("configuration failure");
assert!(failure.configuration);
assert_eq!(failure.session_id.as_deref(), Some("coder-1"));
assert_eq!(failure.detail, "no independent coder");
assert!(configuration_run_failure("coder-1", "red".into(), Some("error")).is_none());
}
async fn poll(
reviewers: Vec<Arc<dyn Reviewer>>,
criteria: &str,
diff: &str,
) -> (Vec<Verdict>, Vec<Unreachable>) {
let futures: Vec<_> = reviewers
.iter()
.map(|r| async move {
let model = r.model().to_string();
match r.review(criteria, diff).await {
Ok(answer) => parse_verdict(&model, &answer),
Err(e) => Err(Unreachable { model, error: e }),
}
})
.collect();
let mut verdicts = Vec::new();
let mut unreachable = Vec::new();
for result in futures::future::join_all(futures).await {
match result {
Ok(v) => verdicts.push(v),
Err(u) => unreachable.push(u),
}
}
(verdicts, unreachable)
}
#[tokio::test]
async fn a_panel_keeps_answers_and_failures_apart() {
let reviewers: Vec<Arc<dyn Reviewer>> = vec![
Arc::new(Fixed {
model: "a",
answer: Ok("PASS looks right".into()),
}),
Arc::new(Fixed {
model: "b",
answer: Ok("FAIL wrong scope".into()),
}),
Arc::new(Fixed {
model: "c",
answer: Err("429 rate limited".into()),
}),
];
let panel_size = reviewers.len();
let (verdicts, unreachable) = poll(reviewers, "criteria", "diff").await;
assert_eq!(verdicts.len(), 2);
assert_eq!(unreachable.len(), 1);
assert_eq!(unreachable[0].model, "c");
assert_eq!(panel_size, 3);
}
#[tokio::test]
async fn an_unreadable_answer_becomes_unreachable_not_a_pass() {
let (verdicts, unreachable) = poll(
vec![Arc::new(Fixed {
model: "a",
answer: Ok("hmm, hard to say".into()),
})],
"c",
"d",
)
.await;
assert!(verdicts.is_empty());
assert_eq!(unreachable.len(), 1);
}
#[test]
fn a_watch_only_target_is_refused_loudly_not_unwrapped() {
let t = HealTarget {
repo: "acme/w".into(),
fix_repo: None,
checkout: None,
label: "self-heal".into(),
base: "main".into(),
};
assert!(checkout_path(&t).is_err());
}
}