use crate::controller::Controller;
use crate::review_settings::{
ReviewCapabilityChoices, ReviewDiscoveryOutcome, ReviewDiscoveryRequest,
};
use crate::session_manager::ManagedSessionHandle;
use crate::utility_llm::{UtilityLlmRuntime, UtilityQuotaClass, classify_quota};
use anyhow::{Context, Result, bail, ensure};
use mj_core::config::ReviewConfig;
use mj_core::review::settings::{
ResolvedReviewSettings, ReviewModelSettings, ReviewProvider, model_matches_family,
model_version_cmp,
};
use std::sync::{
Arc,
atomic::{AtomicBool, Ordering},
};
#[derive(Debug)]
struct Candidate {
id: String,
provider: ReviewProvider,
group: u8,
quota: UtilityQuotaClass,
score: u8,
advertised: bool,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum CatalogVerdict {
Offers,
Lacks,
Unknown,
}
fn catalog_verdict(
offered: Option<&mj_core::profile_capabilities::ProfileCapabilitiesSnapshot>,
key: &str,
model: &str,
) -> CatalogVerdict {
use mj_core::profile_capabilities::CapabilityState;
match offered
.and_then(|snapshot| snapshot.profiles.get(key))
.map(|entry| &entry.choices)
{
Some(CapabilityState::Ready(choices)) => {
if choices.models.iter().any(|choice| choice.value == model) {
CatalogVerdict::Offers
} else {
CatalogVerdict::Lacks
}
}
_ => CatalogVerdict::Unknown,
}
}
fn rank(candidates: &mut [Candidate]) {
candidates.sort_by(|a, b| {
a.group
.cmp(&b.group)
.then_with(|| b.quota.cmp(&a.quota))
.then_with(|| a.provider.cmp(&b.provider))
.then_with(|| b.score.cmp(&a.score))
.then_with(|| a.id.cmp(&b.id))
});
}
pub(crate) async fn resolve(
handle: ManagedSessionHandle,
settings: Option<ReviewConfig>,
cancelled: Arc<AtomicBool>,
offered: Option<mj_core::profile_capabilities::ProfileCapabilitiesSnapshot>,
) -> Result<ResolvedReviewSettings> {
let controller = Arc::new(tokio::task::spawn_blocking(Controller::load).await??);
let settings = settings.unwrap_or_else(|| controller.config.review.clone());
let session = controller
.state
.sessions
.get(handle.session_id())
.context("review session is missing")?;
ensure!(!session.archived, "this session is archived");
let primary = session.last_profile.clone();
let config = controller.config.clone();
let (providers, primary_provider, primary, mut reasons) =
tokio::task::spawn_blocking(move || {
let primary_provider = config
.profiles
.get(&primary)
.and_then(|profile| ReviewProvider::for_profile(profile).ok())
.unwrap_or(ReviewProvider::Other);
let mut providers = std::collections::BTreeMap::new();
let mut reasons = Vec::new();
for (id, profile) in config.enabled_profiles() {
match ReviewProvider::for_profile(profile) {
Ok(provider) => {
providers.insert(id.to_owned(), provider);
}
Err(error) => {
reasons.push(format!("{id}: could not inspect provider: {error:#}"))
}
}
}
(providers, primary_provider, primary, reasons)
})
.await?;
let mut candidates = Vec::new();
if let Some(id) = &settings.profile {
let provider = providers.get(id).with_context(|| {
format!(
"reviewer profile {id:?} is missing, disabled, or unavailable: {}",
reasons.join("; ")
)
})?;
candidates.push(Candidate {
id: id.clone(),
provider: *provider,
group: 0,
quota: UtilityQuotaClass::Unknown,
score: 0,
advertised: true,
});
} else {
let pinned_model = settings.model.is_some();
let mut advertised = std::collections::BTreeSet::new();
let supported = controller
.config
.enabled_profiles()
.filter(|(id, _)| {
providers
.get(*id)
.is_some_and(|p| pinned_model || p.main_policy().is_some())
})
.filter(|(id, profile)| {
let Some(model) = settings.model.as_deref() else {
return true;
};
match catalog_verdict(offered.as_ref(), &profile.capabilities_key(id), model) {
CatalogVerdict::Offers => {
advertised.insert((*id).to_owned());
true
}
CatalogVerdict::Lacks => {
reasons.push(format!("{id}: does not advertise model {model}"));
false
}
CatalogVerdict::Unknown => true,
}
})
.collect::<Vec<_>>();
let quotas = UtilityLlmRuntime::shared()
.quotas(&controller.config, &supported)
.await;
for (id, _) in supported {
let provider = providers[id];
let Some((quota, score)) = quotas
.get(id)
.map(classify_quota)
.unwrap_or(Some((UtilityQuotaClass::Unknown, 0)))
else {
reasons.push(format!("{id}: quota is exhausted"));
continue;
};
candidates.push(Candidate {
id: id.to_owned(),
provider,
group: if id == primary {
2
} else {
u8::from(provider == primary_provider)
},
quota,
score,
advertised: advertised.contains(id),
});
}
rank(&mut candidates);
candidates.sort_by_key(|candidate| !candidate.advertised);
}
for candidate in candidates {
ensure!(
!cancelled.load(Ordering::Acquire),
"review preparation cancelled"
);
let result = resolve_candidate(
controller.clone(),
&handle,
&candidate,
&settings,
&cancelled,
)
.await;
match result {
Ok(mut resolved) => {
resolved.same_provider = candidate.provider == primary_provider;
return Ok(resolved);
}
Err(error) if settings.profile.is_some() => return Err(error),
Err(error) if preempted_by_lifecycle(&format!("{error:#}")) => return Err(error),
Err(error) => {
tracing::warn!(profile = %candidate.id, %error, "Auto reviewer candidate unavailable");
reasons.push(format!("{}: {error:#}", candidate.id));
}
}
}
if reasons.is_empty() {
reasons.push("configure an enabled Codex, Claude, DeepSeek, or Kimi profile, or select a reviewer manually in Settings".into());
}
bail!("No usable Auto reviewer: {}", reasons.join("; "))
}
async fn discover(
controller: Arc<Controller>,
handle: &ManagedSessionHandle,
profile: &str,
model: Option<String>,
cancelled: &Arc<AtomicBool>,
) -> Result<ReviewCapabilityChoices> {
let (progress, _receiver) = tokio::sync::mpsc::unbounded_channel();
let request = ReviewDiscoveryRequest {
profile: profile.into(),
model,
preferred_session: Some(handle.session_id().into()),
};
let result = crate::review_settings::discover_selected_worker(
controller,
handle.session_id().into(),
crate::review_host::next_review_generation().map_err(anyhow::Error::msg)?,
handle.clone(),
&request,
cancelled,
&progress,
)
.await
.map_err(anyhow::Error::msg)?;
match result {
ReviewDiscoveryOutcome::Available {
choices,
cleanup_warning,
} => {
if let Some(warning) = cleanup_warning {
bail!("{warning}");
}
Ok(choices)
}
ReviewDiscoveryOutcome::Unavailable => bail!("review worker is unavailable"),
}
}
fn family_model(choices: &ReviewCapabilityChoices, family: &str) -> Result<String> {
choices
.model_choices
.iter()
.filter(|c| model_matches_family(&c.value, family))
.max_by(|a, b| model_version_cmp(&a.value, &b.value))
.map(|c| c.value.clone())
.with_context(|| format!("no advertised {family} model"))
}
async fn validate_model(
controller: Arc<Controller>,
handle: &ManagedSessionHandle,
profile: &str,
settings: &ReviewModelSettings,
cancelled: &Arc<AtomicBool>,
) -> Result<()> {
let choices = discover(
controller,
handle,
profile,
settings.model.clone(),
cancelled,
)
.await?;
if let Some(model) = &settings.model {
ensure!(
choices.model_choices.iter().any(|c| &c.value == model),
"reviewer {profile} does not advertise model {model}"
);
}
if let Some(effort) = &settings.effort {
ensure!(
choices.effort_capabilities_discovered
&& choices.effort_choices.iter().any(|c| &c.value == effort),
"reviewer {profile} does not advertise effort {effort} for {}",
settings.model.as_deref().unwrap_or("its default model")
);
}
Ok(())
}
fn select_models(
provider: ReviewProvider,
settings: &ReviewConfig,
catalog: &ReviewCapabilityChoices,
) -> Result<ReviewModelSettings> {
let automatic = settings.profile.is_none();
let main = if automatic {
let policy = provider.main_policy();
let model = match &settings.model {
Some(model) => {
ensure!(
catalog.model_choices.iter().any(|c| &c.value == model),
"does not advertise model {model}"
);
model.clone()
}
None => family_model(catalog, policy.context("provider is manual-only")?.0)?,
};
ReviewModelSettings {
model: Some(model),
effort: settings
.effort
.clone()
.or_else(|| policy.map(|(_, effort)| effort.into())),
fast_mode: false,
}
} else {
ReviewModelSettings {
model: settings.model.clone(),
effort: settings.effort.clone(),
fast_mode: false,
}
};
Ok(main)
}
async fn resolve_candidate(
controller: Arc<Controller>,
handle: &ManagedSessionHandle,
candidate: &Candidate,
settings: &ReviewConfig,
cancelled: &Arc<AtomicBool>,
) -> Result<ResolvedReviewSettings> {
let automatic = settings.profile.is_none();
let catalog = discover(controller.clone(), handle, &candidate.id, None, cancelled).await?;
let main = select_models(candidate.provider, settings, &catalog)?;
if main.model.is_some() || main.effort.is_some() {
validate_model(controller.clone(), handle, &candidate.id, &main, cancelled).await?;
}
Ok(ResolvedReviewSettings {
profile: candidate.id.clone(),
generation: crate::review_host::next_review_generation().map_err(anyhow::Error::msg)?,
main,
automatic,
same_provider: false,
})
}
pub(crate) fn preempted_by_lifecycle(reason: &str) -> bool {
reason.contains("session is reserved for a lifecycle operation")
|| reason.contains("cancelled for session lifecycle change")
|| reason.contains("worker is reserved for checkpoint or replacement")
}
#[cfg(test)]
mod tests {
use super::*;
fn candidate(
id: &str,
provider: ReviewProvider,
group: u8,
quota: UtilityQuotaClass,
score: u8,
) -> Candidate {
Candidate {
id: id.into(),
provider,
group,
quota,
score,
advertised: false,
}
}
fn snapshot(
entries: &[(
&str,
mj_core::profile_capabilities::CapabilityState<Vec<&str>>,
)],
) -> mj_core::profile_capabilities::ProfileCapabilitiesSnapshot {
use mj_core::profile_capabilities::{CapabilityState, ProfileCapabilities};
mj_core::profile_capabilities::ProfileCapabilitiesSnapshot {
profiles: entries
.iter()
.map(|(key, state)| {
let choices = match state {
CapabilityState::Pending => CapabilityState::Pending,
CapabilityState::Failed(error) => CapabilityState::Failed(error.clone()),
CapabilityState::Ready(models) => {
CapabilityState::Ready(mj_core::worker_launch::ProfileConfig {
model: None,
models: catalog(models).model_choices,
efforts: Vec::new(),
observed_at: 0,
})
}
};
(
(*key).to_owned(),
ProfileCapabilities {
choices,
efforts: Default::default(),
},
)
})
.collect(),
}
}
#[test]
fn a_pinned_model_rules_profiles_in_or_out_from_the_catalog_alone() {
use mj_core::profile_capabilities::CapabilityState;
let offered = snapshot(&[
(
"codex3",
CapabilityState::Ready(vec!["gpt-6-luna", "gpt-6.1-sol"]),
),
("kimi", CapabilityState::Ready(vec!["k3", "k3-256k"])),
(
"muse",
CapabilityState::Failed("discovery harness stopped".into()),
),
("claude", CapabilityState::Pending),
]);
let verdict = |key| catalog_verdict(Some(&offered), key, "gpt-6-luna");
assert_eq!(verdict("codex3"), CatalogVerdict::Offers);
assert_eq!(verdict("kimi"), CatalogVerdict::Lacks);
assert_eq!(verdict("muse"), CatalogVerdict::Unknown);
assert_eq!(verdict("claude"), CatalogVerdict::Unknown);
assert_eq!(verdict("glm"), CatalogVerdict::Unknown);
assert_eq!(
catalog_verdict(None, "codex3", "gpt-6-luna"),
CatalogVerdict::Unknown
);
}
#[test]
fn confirmed_profiles_are_tried_before_unknown_ones_whatever_their_rank() {
let mut candidates = vec![
candidate(
"muse",
ReviewProvider::Other,
0,
UtilityQuotaClass::Unknown,
0,
),
Candidate {
advertised: true,
..candidate(
"codex3",
ReviewProvider::Codex,
1,
UtilityQuotaClass::Unknown,
0,
)
},
candidate(
"glm",
ReviewProvider::Other,
0,
UtilityQuotaClass::Unknown,
0,
),
];
rank(&mut candidates);
candidates.sort_by_key(|candidate| !candidate.advertised);
assert_eq!(candidates[0].id, "codex3");
}
#[test]
fn review_selection_uses_one_model_and_effort_from_an_acp_config_fixture() {
let options: Vec<agent_client_protocol::schema::v1::SessionConfigOption> =
serde_json::from_str(include_str!(
"../tests/fixtures/review-acp-config-options.json"
))
.expect("ACP v1 session config options");
let catalog = ReviewCapabilityChoices {
model_choices: mj_core::acp::session_config_choices(&options, "model"),
effort_choices: mj_core::acp::session_config_choices(&options, "effort"),
effort_capabilities_discovered: true,
};
let reviewer = select_models(ReviewProvider::Codex, &ReviewConfig::default(), &catalog)
.expect("select the advertised Codex review model");
assert_eq!(reviewer.model.as_deref(), Some("gpt-6.10-astra"));
assert_eq!(reviewer.effort.as_deref(), Some("medium"));
assert!(!reviewer.fast_mode);
}
fn catalog(ids: &[&str]) -> ReviewCapabilityChoices {
ReviewCapabilityChoices {
model_choices: ids
.iter()
.map(|id| mj_core::acp::SessionConfigChoice {
value: (*id).into(),
name: (*id).into(),
description: None,
})
.collect(),
..Default::default()
}
}
}
#[cfg(test)]
mod preemption_tests {
use super::preempted_by_lifecycle;
#[test]
fn a_worker_reserved_for_checkpoint_is_a_lifecycle_preemption() {
let refusal = "relay 2.28.0 could not perform reviewer_start: relay rejected request \
(InvalidState): worker is reserved for checkpoint or replacement; \
reviewer work was not admitted";
assert!(preempted_by_lifecycle(refusal));
assert!(preempted_by_lifecycle(
"session is reserved for a lifecycle operation"
));
assert!(!preempted_by_lifecycle(
"does not advertise model gpt-6-luna"
));
}
}