use crate::cli::output::{OutputConfig, OutputFormat};
use crate::cli::{ModelRelayCommands, ModelRelayToggleArgs};
use crate::error::{OlError, ERR_MODEL_RELAY_FINDING_NOT_FOUND};
pub fn run(cmd: &ModelRelayCommands, output: &OutputConfig) -> Result<(), OlError> {
match cmd {
ModelRelayCommands::Status => status(output),
ModelRelayCommands::Enable(args) => toggle(true, args, output),
ModelRelayCommands::Disable(args) => toggle(false, args, output),
ModelRelayCommands::Explain { finding_id } => explain(finding_id, output),
}
}
fn toggle(enable: bool, args: &ModelRelayToggleArgs, output: &OutputConfig) -> Result<(), OlError> {
use std::io::{BufRead, IsTerminal, Write};
let config_path = crate::config::openlatch_dir().join("config.toml");
if !config_path.exists() {
crate::config::ensure_config(crate::config::Config::defaults().port)?;
}
let before = crate::config::Config::load(None, None, false)
.map(|c| c.model_relay.enabled)
.unwrap_or(true);
let verb = if enable { "enabled" } else { "disabled" };
crate::cli::header::print(
output,
&[
"system",
"model-relay",
if enable { "enable" } else { "disable" },
],
);
if before == enable {
output.print_substep(&format!("Model relay already {verb} in config"));
} else {
crate::config::persist_model_relay_enabled(&config_path, enable)?;
output.print_step(&format!("Model relay {verb} in {}", config_path.display()));
}
let daemon_up = crate::cli::commands::lifecycle::read_pid_file()
.map(crate::cli::commands::lifecycle::is_process_alive)
.unwrap_or(false);
if !daemon_up {
output.print_step("No daemon running — the change applies at the next start");
emit_toggle_json(enable, true, false, true, output);
return Ok(());
}
let interactive =
std::io::stdin().is_terminal() && output.format == OutputFormat::Human && !output.quiet;
let restart = if args.yes {
true
} else if args.no_restart || !interactive {
false
} else {
eprint!("Restart the daemon now to apply? [y/N] ");
let _ = std::io::stderr().flush();
let mut answer = String::new();
let _ = std::io::stdin().lock().read_line(&mut answer);
matches!(answer.trim().to_ascii_lowercase().as_str(), "y" | "yes")
};
if !restart {
output.print_substep(
"Config updated — not in effect until the daemon restarts (run `openlatch restart`)",
);
emit_toggle_json(enable, true, false, false, output);
crate::cli::report::record_exit_code(crate::cli::report::EXIT_DEGRADED);
return Ok(());
}
crate::cli::commands::lifecycle::run_restart(output)?;
let cfg = crate::config::Config::load(None, None, false)?;
let in_effect = model_relay_rows(&cfg).iter().all(|row| match row.state {
ModelRelayState::Disabled => !enable,
ModelRelayState::Wired | ModelRelayState::Isolated => enable,
_ => false,
});
if in_effect {
output.print_step(&format!("Model relay {verb} and in effect"));
} else {
output.print_substep(
"Daemon restarted, but the model relay is not in the requested state — run \
`openlatch doctor` for the reason",
);
crate::cli::report::record_exit_code(crate::cli::report::EXIT_DEGRADED);
}
emit_toggle_json(enable, true, true, in_effect, output);
Ok(())
}
fn emit_toggle_json(
enabled: bool,
config_written: bool,
restarted: bool,
in_effect: bool,
output: &OutputConfig,
) {
if output.format != OutputFormat::Json {
return;
}
output.print_json(&serde_json::json!({
"enabled": enabled,
"config_written": config_written,
"restarted": restarted,
"in_effect": in_effect,
"exit_code": if in_effect { 0 } else { crate::cli::report::EXIT_DEGRADED },
}));
}
#[derive(Debug, PartialEq)]
pub(crate) enum PortOwnership {
Owned,
Foreign,
Unreachable,
}
const CONNECT_TIMEOUT: std::time::Duration = std::time::Duration::from_millis(100);
const PROBE_TIMEOUT: std::time::Duration = std::time::Duration::from_millis(500);
fn get_admin_status(port: u16) -> Option<reqwest::blocking::Response> {
let url = format!("http://127.0.0.1:{port}/admin/model-relay/status");
crate::egress::blocking_client_builder()
.timeout(PROBE_TIMEOUT)
.build()
.ok()?
.get(&url)
.send()
.ok()
}
pub(crate) fn verify_port_ownership(port: u16) -> PortOwnership {
match relay_status_on(port) {
Ok(_) => PortOwnership::Owned,
Err(ownership) => ownership,
}
}
pub(crate) fn verify_endpoint_ownership(port: u16, key: &str) -> PortOwnership {
match relay_status_on(port) {
Ok(v) if v["endpoint"]["key"].as_str() == Some(key) => PortOwnership::Owned,
Ok(_) => PortOwnership::Foreign,
Err(ownership) => ownership,
}
}
fn relay_status_on(port: u16) -> Result<serde_json::Value, PortOwnership> {
let addr = std::net::SocketAddr::from(([127, 0, 0, 1], port));
if std::net::TcpStream::connect_timeout(&addr, CONNECT_TIMEOUT).is_err() {
return Err(PortOwnership::Unreachable);
}
let Some(resp) = get_admin_status(port) else {
return Err(PortOwnership::Foreign);
};
if !resp.status().is_success() {
return Err(PortOwnership::Foreign);
}
match resp.json::<serde_json::Value>() {
Ok(v) if v.get("status").is_some() && v.get("upstream").is_some() => Ok(v),
_ => Err(PortOwnership::Foreign),
}
}
pub fn status(output: &OutputConfig) -> Result<(), OlError> {
let cfg = crate::config::Config::load(None, None, false)?;
let port = cfg.model_relay.port;
let rows = model_relay_rows(&cfg);
let probe = probe_model_relay(port);
let (endpoints, endpoint_checks) = endpoint_status(
&endpoint_rows_for(
&cfg,
&crate::hooks::detect_agents(),
&verify_endpoint_ownership,
),
probe.as_ref(),
);
let severity =
model_relay_severity(&worst_row(&rows).state).max(endpoint_severity(&endpoint_checks));
let exit = match severity {
0 => 0,
2 => 1,
_ => crate::cli::report::EXIT_DEGRADED,
};
crate::cli::report::record_exit_code(exit);
if output.format == OutputFormat::Json {
output.print_json(&serde_json::json!({
"port": port,
"state": worst_row(&rows).state.label(),
"upstream": crate::model_relay::wire_format::WireFormat::ALL
.iter()
.map(|f| (f.as_str().to_string(), serde_json::json!(cfg.model_relay.upstream_for(*f))))
.collect::<serde_json::Map<_, _>>(),
"classification": format!("{:?}", worst_row(&rows).state),
"enabled": cfg.model_relay.enabled,
"owns_agent_wiring": cfg.model_relay.owns_agent_wiring(),
"wired_to": worst_row(&rows).wired,
"agents": rows
.iter()
.filter(|r| !r.agent.is_empty())
.map(agent_row_json)
.collect::<Vec<_>>(),
"endpoints": endpoints,
"up": probe.is_some(),
"detail": probe,
"exit_code": exit,
}));
return Ok(());
}
crate::cli::header::print(output, &["model-relay status"]);
let multi = rows.len() > 1;
for row in &rows {
let wired = row.wired.clone();
let (line, remedy): (String, Option<String>) = match &row.state {
ModelRelayState::Disabled => (
"Model relay: disabled in config — model calls bypass OpenLatch".to_string(),
Some("Run `openlatch system model-relay enable` to turn it back on.".to_string()),
),
ModelRelayState::Isolated => (
format!("Model relay: isolated instance on port {port}"),
Some(format!(
"This instance does not touch the machine-global agent config. Route a \
session through it with:\n {}",
row.route_hint
)),
),
ModelRelayState::Wired => (
format!(
"Model relay: up on port {port}, agent wired to {}",
wired.as_deref().unwrap_or("it")
),
None,
),
ModelRelayState::WiredButDown => (
format!("Model relay: agent is wired to 127.0.0.1:{port} but nothing is listening"),
Some(
"Model calls fail with ECONNREFUSED. Run `openlatch start` to bring the \
listener up, or `openlatch stop` to clear the wiring and go direct."
.to_string(),
),
),
ModelRelayState::WiredToForeign => (
format!("Model relay: 127.0.0.1:{port} is held by a process that is NOT OpenLatch"),
Some(format!(
"The agent is wired to it, so your provider API key is going to that process. \
Identify it (lsof -i :{port}), stop it, then run `openlatch restart`."
)),
),
ModelRelayState::PreflightFailed(why) => (
format!("Model relay: up on port {port}, preflight FAILED — {why}"),
Some(
"The agent was left unwired on purpose: model calls go direct and keep \
working, but nothing is captured. Fix reachability to the provider; the \
daemon re-wires itself as soon as the check passes."
.to_string(),
),
),
ModelRelayState::PreflightPending => (
format!("Model relay: up on port {port}, preflight still running"),
Some("The agent is wired once it passes. Re-run this in a moment.".to_string()),
),
ModelRelayState::UpUnwired => (
format!("Model relay: up on port {port} but the agent is not wired to it"),
Some("Run `openlatch restart` to re-wire.".to_string()),
),
ModelRelayState::Down => (
format!("Model relay: enabled in config, nothing listening on port {port}"),
Some("Run `openlatch start`.".to_string()),
),
ModelRelayState::ForeignIdle => (
format!("Model relay: 127.0.0.1:{port} is held by another process"),
Some(format!(
"The agent is not wired to it, but the next `openlatch start` will refuse to \
bind. Identify it with `lsof -i :{port}`."
)),
),
};
if multi {
eprintln!(" [{}]", row.display_name);
}
eprintln!(" {line}");
if let Some(remedy) = remedy {
eprintln!(" {remedy}");
}
}
for check in &endpoint_checks {
eprintln!(" {}", check.headline);
if check.state != crate::cli::report::State::Ok {
if let Some(remedy) = &check.remedy {
eprintln!(" {remedy}");
}
}
}
if let Some(v) = probe.as_ref() {
if let Some(up) = v.get("upstream").and_then(|x| x.as_object()) {
for (fmt, base) in up {
if let Some(base) = base.as_str() {
eprintln!(" Upstream ({fmt}): {base}");
}
}
}
if let Some(base) = v.get("upstream_chatgpt").and_then(|x| x.as_str()) {
eprintln!(" Upstream (openai-responses, ChatGPT plan): {base}");
}
if let Some(f) = v.get("pass_through_failures").and_then(|x| x.as_u64()) {
eprintln!(" Pass-through failures: {f}");
}
}
Ok(())
}
pub fn explain(finding_id: &str, output: &OutputConfig) -> Result<(), OlError> {
let record = crate::model_relay::retention::load(finding_id).ok_or_else(|| {
OlError::new(
ERR_MODEL_RELAY_FINDING_NOT_FOUND,
format!("no local churn finding '{finding_id}'"),
)
.with_suggestion(
"Findings resolve only on the host that produced them, and expire from the bounded \
local store. Check the id from the `ai.openlatch.prefix.finding_id` field.",
)
})?;
if output.format == OutputFormat::Json {
output.print_json(&serde_json::json!({
"finding_id": record.finding_id,
"captured_at": record.captured_at,
"churn_layer": record.churn_layer,
"churn_class": record.churn_class,
"divergence_offset": record.divergence_offset,
"churn_byte_len": record.churn_byte_len,
"churn_block_index": record.churn_block_index,
"block": record.block,
}));
} else {
crate::cli::header::print(output, &["model-relay explain"]);
eprintln!(" finding : {}", record.finding_id);
eprintln!(" captured : {}", record.captured_at);
eprintln!(" layer : {}", record.churn_layer);
eprintln!(" class : {}", record.churn_class);
eprintln!(
" offset/len : {} / {} (block #{})",
record.divergence_offset, record.churn_byte_len, record.churn_block_index
);
eprintln!(" block (local, never emitted):");
println!("{}", record.block);
}
Ok(())
}
pub fn probe_model_relay(port: u16) -> Option<serde_json::Value> {
let resp = get_admin_status(port)?;
if !resp.status().is_success() {
return None;
}
resp.json().ok()
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) enum ModelRelayState {
Disabled,
Isolated,
Wired,
WiredButDown,
WiredToForeign,
PreflightFailed(String),
PreflightPending,
UpUnwired,
Down,
ForeignIdle,
}
impl ModelRelayState {
pub(crate) fn label(&self) -> &'static str {
match self {
ModelRelayState::Disabled => "disabled",
ModelRelayState::Isolated => "isolated",
ModelRelayState::Wired => "up",
ModelRelayState::WiredButDown => "down",
ModelRelayState::WiredToForeign => "failed",
ModelRelayState::PreflightFailed(_) => "preflight-failed",
ModelRelayState::PreflightPending => "preflight-pending",
ModelRelayState::UpUnwired => "unwired",
ModelRelayState::Down => "down",
ModelRelayState::ForeignIdle => "down",
}
}
}
pub(crate) fn agent_row_json(row: &AgentRelay) -> serde_json::Value {
serde_json::json!({
"agent": row.agent,
"state": row.state.label(),
"classification": format!("{:?}", row.state),
"wired_to": row.wired,
"wireformat": row.wire_format.map(|f| f.as_str()),
})
}
pub(crate) struct AgentRelay {
pub agent: &'static str,
pub display_name: &'static str,
pub wired: Option<String>,
pub state: ModelRelayState,
pub wire_format: Option<crate::model_relay::wire_format::WireFormat>,
pub route_hint: String,
}
pub(crate) fn model_relay_rows(cfg: &crate::config::Config) -> Vec<AgentRelay> {
model_relay_rows_for(cfg, &crate::hooks::detect_agents())
}
pub(crate) fn model_relay_rows_for(
cfg: &crate::config::Config,
agents: &[crate::hooks::DetectedAgent],
) -> Vec<AgentRelay> {
let port = cfg.model_relay.port;
let mut rows: Vec<AgentRelay> = agents
.iter()
.filter_map(|a| {
let wire_format = Some(a.binding.model_relay_wiring()?.wire_format);
let wired = read_agent_wiring(&*a.binding);
let state = classify_model_relay(cfg, a.agent_type(), wired.as_deref());
Some(AgentRelay {
agent: a.agent_type(),
display_name: a.binding.display_name(),
wired,
state,
wire_format,
route_hint: crate::hooks::isolated_wiring_hint(&*a.binding, port),
})
})
.collect();
if rows.is_empty() {
rows.push(AgentRelay {
agent: "",
display_name: "The agent",
wired: None,
state: classify_model_relay(cfg, "", None),
wire_format: None,
route_hint: "the base URL your agent reads".to_string(),
});
}
rows
}
pub(crate) fn model_relay_severity(state: &ModelRelayState) -> u8 {
match state {
ModelRelayState::Wired | ModelRelayState::Isolated => 0,
ModelRelayState::Disabled
| ModelRelayState::PreflightPending
| ModelRelayState::UpUnwired
| ModelRelayState::Down => 1,
ModelRelayState::WiredButDown
| ModelRelayState::WiredToForeign
| ModelRelayState::PreflightFailed(_)
| ModelRelayState::ForeignIdle => 2,
}
}
pub(crate) fn worst_row(rows: &[AgentRelay]) -> &AgentRelay {
rows.iter()
.max_by_key(|r| model_relay_severity(&r.state))
.expect("model_relay_rows never returns an empty list")
}
pub(crate) fn classify_model_relay(
cfg: &crate::config::Config,
agent: &'static str,
wired: Option<&str>,
) -> ModelRelayState {
if !cfg.model_relay.enabled {
return ModelRelayState::Disabled;
}
if !cfg.model_relay.owns_agent_wiring() {
return ModelRelayState::Isolated;
}
let port = cfg.model_relay.port;
match (wired, verify_port_ownership(port)) {
(Some(_), PortOwnership::Owned) => ModelRelayState::Wired,
(Some(_), PortOwnership::Unreachable) => ModelRelayState::WiredButDown,
(Some(_), PortOwnership::Foreign) => ModelRelayState::WiredToForeign,
(None, PortOwnership::Owned) => {
let live = probe_model_relay(port);
match live
.as_ref()
.and_then(|v| v.get("preflight"))
.and_then(|v| v.get(agent))
.and_then(|v| v.as_str())
{
Some("failed") => ModelRelayState::PreflightFailed(
live.as_ref()
.and_then(|v| v.get("preflight_error"))
.and_then(|v| v.get(agent))
.and_then(|v| v.as_str())
.unwrap_or("no round trip to the provider completed")
.to_string(),
),
Some("pending") => ModelRelayState::PreflightPending,
_ => ModelRelayState::UpUnwired,
}
}
(None, PortOwnership::Unreachable) => ModelRelayState::Down,
(None, PortOwnership::Foreign) => ModelRelayState::ForeignIdle,
}
}
pub(crate) fn read_agent_wiring(
binding: &dyn crate::hooks::binding::AgentBinding,
) -> Option<String> {
use crate::hooks::binding::EndpointConvention;
match binding.model_relay_wiring()?.endpoint {
EndpointConvention::EnvVars { .. } => {
read_model_relay_base_url(&binding.hook_config_path())
}
EndpointConvention::TomlProvider { provider_name, .. } => {
crate::hooks::codex_cli::read_provider_base_url(
&crate::hooks::codex_cli::config_toml_path(&binding.config_dir()),
provider_name,
)
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct EndpointRow {
pub agent: &'static str,
pub key: String,
pub label: String,
pub file: Option<std::path::PathBuf>,
pub port: Option<u16>,
pub state: EndpointState,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) enum EndpointState {
Active,
NextStart,
Unwired,
Down,
Foreign,
Verdict {
code: &'static str,
detail: String,
},
Uncovered(crate::hooks::cline_providers::UncoveredReason),
StateFile(String),
}
pub(crate) fn endpoint_rows_for(
cfg: &crate::config::Config,
agents: &[crate::hooks::DetectedAgent],
owner: &dyn Fn(u16, &str) -> PortOwnership,
) -> Vec<EndpointRow> {
if !cfg.model_relay.enabled {
return Vec::new();
}
let mut status: Option<Option<serde_json::Value>> = None;
let mut rows = Vec::new();
for agent in agents {
let Some(endpoints) = agent.binding.provider_endpoints() else {
continue;
};
if !crate::daemon::owns_wiring_for(cfg, &*agent.binding) {
continue;
}
let status = status.get_or_insert_with(|| probe_model_relay(cfg.model_relay.port));
rows.extend(rows_for_agent(cfg, endpoints, status.as_ref(), owner));
}
rows
}
fn rows_for_agent(
cfg: &crate::config::Config,
endpoints: &dyn crate::hooks::provider_endpoints::ProviderEndpoints,
status: Option<&serde_json::Value>,
owner: &dyn Fn(u16, &str) -> PortOwnership,
) -> Vec<EndpointRow> {
use crate::hooks::cline_providers::{
decide, DecideCtx, Decision, SlotId, SlotObservation, UncoveredReason,
};
use crate::model_relay::endpoints::RelayPorts;
let agent = endpoints.agent_type();
let row = |key: String, label: String, file, port, state| EndpointRow {
agent,
key,
label,
file,
port,
state,
};
let records: std::collections::BTreeMap<
String,
crate::hooks::model_relay_endpoints::EndpointRecord,
> = match crate::hooks::model_relay_endpoints::endpoint_records(endpoints.record_prefix()) {
Ok(records) => records.into_iter().collect(),
Err(e) => {
return vec![row(
endpoints.record_prefix().to_string(),
"provider endpoint records".into(),
None,
None,
EndpointState::StateFile(e.message),
)]
}
};
let daemon_up = status.is_some();
let verdict = |key: &str| {
let v = status?.get("endpoint_verdicts")?.get(key)?;
let code = v.get("code")?.as_str()?;
let code = crate::error::ERR_MODEL_RELAY_CODES
.iter()
.copied()
.find(|c| *c == code)?;
Some(EndpointState::Verdict {
code,
detail: v
.get("detail")
.and_then(|d| d.as_str())
.unwrap_or_default()
.to_string(),
})
};
let last_request = |key: &str| {
status
.and_then(|s| s.get("endpoints"))
.and_then(|e| e.as_array())
.and_then(|all| all.iter().find(|e| e["key"].as_str() == Some(key)))
.and_then(|e| e["last_request_unix"].as_u64())
};
let label = |obs: &SlotObservation| {
let setting = match &obs.slot {
SlotId::GlobalState { key, .. } => key.clone(),
SlotId::ProvidersJson { .. } => "providers.json".to_string(),
};
let ids = match (&obs.slot, obs.row) {
(SlotId::ProvidersJson { id }, _) => id.clone(),
(_, Some(row)) => row.ids.join(", "),
(_, None) => obs.selected_ids.join(", "),
};
format!("{ids} ({setting})")
};
let recorded = records.keys().cloned().collect();
let observation = endpoints.observe(&recorded);
let ctx = DecideCtx {
ports: RelayPorts {
daemon: cfg.port,
main: cfg.model_relay.port,
},
started_at: u64::MAX,
};
let mut rows = Vec::new();
for problem in &observation.problems {
let (path, why) = match problem {
crate::hooks::cline_providers::FileProblem::TooLarge { path, size } => (
path.clone(),
format!(
"{} is {size} bytes, over the limit OpenLatch reads",
crate::core::path_compat::display_path(path)
),
),
crate::hooks::cline_providers::FileProblem::Unparseable { path } => (
path.clone(),
format!(
"{} is not readable as JSON",
crate::core::path_compat::display_path(path)
),
),
};
rows.push(row(
format!("{}file", endpoints.record_prefix()),
"provider settings".into(),
Some(path),
None,
EndpointState::StateFile(why),
));
}
let mut covered: std::collections::BTreeSet<String> = std::collections::BTreeSet::new();
for obs in &observation.slots {
let key = obs.slot.record_key();
covered.extend(obs.selected_ids.iter().cloned());
if let SlotId::ProvidersJson { id } = &obs.slot {
covered.insert(id.clone());
}
let rec = records.get(&key);
let port = rec.map(|r| r.port);
let state = if let Some(v) = verdict(&key) {
Some(v)
} else {
match decide(obs, rec, &ctx) {
Decision::Keep => {
use crate::hooks::model_relay_endpoints::Proof;
let rec = rec.filter(|r| r.is_live());
rec.map(|r| match owner(r.port, &key) {
PortOwnership::Owned
if r.served_since_written(last_request(&key))
|| r.proven_by == Some(Proof::Traffic) =>
{
EndpointState::Active
}
PortOwnership::Owned if r.misconfigured_event.is_some() => {
EndpointState::Verdict {
code: crate::error::ERR_MODEL_RELAY_MISCONFIGURED,
detail: "the editor uses this provider and no request reaches \
its endpoint"
.to_string(),
}
}
PortOwnership::Owned => match obs.traffic_only() {
Some(reason) => EndpointState::Uncovered(reason),
None if r.proven_at.is_some() => EndpointState::Active,
None => EndpointState::NextStart,
},
PortOwnership::Unreachable => EndpointState::Down,
PortOwnership::Foreign => EndpointState::Foreign,
})
}
Decision::Wire { .. } | Decision::Reapply | Decision::NewPrior { .. } => {
Some(if daemon_up {
EndpointState::Unwired
} else if rec.is_some_and(|r| r.is_live()) {
EndpointState::Down
} else {
EndpointState::Unwired
})
}
Decision::Uncovered(reason) => Some(EndpointState::Uncovered(reason)),
Decision::Release { .. } | Decision::Skip => None,
}
};
if let Some(state) = state {
rows.push(row(key, label(obs), Some(obs.file.clone()), port, state));
}
}
let mut reported = std::collections::BTreeSet::new();
for selection in &observation.selections {
if covered.contains(&selection.id) || !reported.insert(selection.id.clone()) {
continue;
}
let reason = match crate::hooks::cline_providers::providers_json_row(&selection.id) {
Some(row) => row.never.unwrap_or(UncoveredReason::NoSetting),
None => UncoveredReason::DefaultUnknown,
};
rows.push(row(
format!("{}selected:{}", endpoints.record_prefix(), selection.id),
selection.id.clone(),
None,
None,
EndpointState::Uncovered(reason),
));
}
rows
}
pub(crate) fn endpoint_status(
rows: &[EndpointRow],
status: Option<&serde_json::Value>,
) -> (Vec<serde_json::Value>, Vec<crate::cli::report::Check>) {
let mut report = crate::cli::report::Report::new();
crate::cli::commands::doctor::check_provider_endpoints(rows, &mut report);
let records: std::collections::BTreeMap<
String,
crate::hooks::model_relay_endpoints::EndpointRecord,
> = crate::hooks::model_relay_endpoints::endpoint_records("")
.map(|all| all.into_iter().collect())
.unwrap_or_default();
let live = |key: &str| {
status
.and_then(|s| s.get("endpoints"))
.and_then(|e| e.as_array())
.and_then(|all| all.iter().find(|e| e["key"].as_str() == Some(key)))
};
let entries = rows
.iter()
.zip(report.checks())
.map(|(row, check)| {
let rec = records.get(&row.key);
let live = live(&row.key);
serde_json::json!({
"agent": row.agent,
"key": row.key,
"label": row.label,
"port": row.port,
"state": check.state.key(),
"code": check.code,
"headline": check.headline,
"remedy": check.remedy,
"origin": rec.map(|r| crate::core::egress::credentials::mask_userinfo(&r.origin)),
"proven_at": rec.and_then(|r| r.proven_at),
"proven_by": rec.and_then(|r| r.proven_by),
"requests": live.and_then(|e| e["requests"].as_u64()),
"contested": live
.and_then(|e| e["contested"].as_bool())
.unwrap_or(false),
})
})
.collect();
(entries, report.checks().to_vec())
}
pub(crate) fn endpoint_severity(checks: &[crate::cli::report::Check]) -> u8 {
checks
.iter()
.map(|c| {
if c.state.is_failure() {
2
} else if c.state.is_warning() {
1
} else {
0
}
})
.max()
.unwrap_or(0)
}
pub(crate) fn read_model_relay_base_url(settings_path: &std::path::Path) -> Option<String> {
let raw = std::fs::read_to_string(settings_path).ok()?;
let parsed = crate::hooks::jsonc::parse_settings_value(&raw).ok()?;
let url = parsed
.get("env")?
.get("ANTHROPIC_BASE_URL")?
.as_str()?
.to_string();
reqwest::Url::parse(url.trim())
.ok()
.filter(|u| u.host_str() == Some("127.0.0.1"))
.map(|_| url)
}
#[cfg(test)]
mod endpoint_row_tests {
use super::{endpoint_rows_for, EndpointState, PortOwnership};
use crate::hooks::cline_providers::{
state_lanes_from, ClineProviderEndpoints, UncoveredReason,
};
use crate::hooks::model_relay_endpoints::{self, EndpointRecord, SlotState, SlotValue};
use std::path::{Path, PathBuf};
struct Host {
_env: crate::hooks::cline::EnvOverride,
_lock: std::sync::MutexGuard<'static, ()>,
_dir: tempfile::TempDir,
root: PathBuf,
}
fn host() -> Host {
let lock = crate::config::OPENLATCH_DIR_ENV_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
let dir = tempfile::tempdir().expect("tempdir");
let root = dir.path().to_path_buf();
let env = crate::hooks::cline::EnvOverride::apply([(
"OPENLATCH_DIR",
Some(root.join("openlatch").into_os_string()),
)]);
Host {
_env: env,
_lock: lock,
_dir: dir,
root,
}
}
#[test]
fn host_never_lets_a_racing_host_see_its_openlatch_dir_stomped() {
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc;
let corrupted = Arc::new(AtomicBool::new(false));
let stop = Arc::new(AtomicBool::new(false));
let worker = |corrupted: Arc<AtomicBool>, stop: Arc<AtomicBool>| {
std::thread::spawn(move || {
for _ in 0..500 {
if stop.load(Ordering::Relaxed) {
return;
}
let h = host();
let expected = h.root.join("openlatch").into_os_string();
for _ in 0..200 {
if std::env::var_os("OPENLATCH_DIR").as_deref() != Some(&expected) {
corrupted.store(true, Ordering::Relaxed);
stop.store(true, Ordering::Relaxed);
break;
}
std::thread::yield_now();
}
drop(h);
}
})
};
let a = worker(corrupted.clone(), stop.clone());
let b = worker(corrupted.clone(), stop.clone());
a.join().expect("thread a");
b.join().expect("thread b");
assert!(
!corrupted.load(Ordering::Relaxed),
"a racing host() observed its OPENLATCH_DIR stomped by another \
guard's delayed restore — the lock released before the env was \
put back"
);
}
fn config() -> crate::config::Config {
let mut cfg = crate::config::Config::defaults();
cfg.model_relay.port = crate::core::egress::test_support::dead_port();
cfg.model_relay.own_agent_wiring = Some(true);
cfg
}
fn agent(root: &Path, state: &str, providers: Option<&str>) -> crate::hooks::DetectedAgent {
let gs = root.join("data").join("globalState.json");
std::fs::create_dir_all(gs.parent().expect("parent")).expect("mkdir");
std::fs::write(&gs, state).expect("write");
let pj = providers.map(|body| {
let path = root.join("data").join("settings").join("providers.json");
std::fs::create_dir_all(path.parent().expect("parent")).expect("mkdir");
std::fs::write(&path, body).expect("write");
path
});
let endpoints: &'static ClineProviderEndpoints = Box::leak(Box::new(
ClineProviderEndpoints::at(state_lanes_from(Some(gs.clone()), Some(gs)), pj),
));
crate::hooks::DetectedAgent {
kind: crate::hooks::AgentKind::Cline,
binding: std::sync::Arc::new(crate::hooks::binding::test_support::FakeBinding {
agent_type: "cline",
provider_endpoints: Some(endpoints),
..Default::default()
}),
}
}
fn wired(root: &Path, port: u16, proven: bool) -> EndpointRecord {
EndpointRecord {
port,
origin: "https://generativelanguage.googleapis.com/".into(),
prior: SlotValue::Absent,
last_written: Some(format!("http://127.0.0.1:{port}")),
file: root.join("data").join("globalState.json"),
state: SlotState::Wired,
released_by: None,
changed_at: 1,
proven_at: proven.then_some(2),
family: None,
pending_event: None,
proven_by: None,
misconfigured_event: None,
}
}
#[test]
fn endpoint_rows_are_empty_without_records_or_configuration() {
let host = host();
let agents = [agent(&host.root, "{}", None)];
let rows = endpoint_rows_for(&config(), &agents, &|_, _| PortOwnership::Owned);
assert!(rows.is_empty(), "{rows:?}");
}
#[test]
fn each_slot_is_classified_from_its_record_its_port_and_its_file() {
let host = host();
let cfg = config();
let block = crate::model_relay::endpoints::endpoint_port_block(cfg.model_relay.port);
let (p1, p2, p3, p4) = (
*block.start(),
block.start() + 1,
block.start() + 2,
block.start() + 3,
);
let state = format!(
r#"{{"actModeApiProvider":"ollama","planModeApiProvider":"bedrock",
"geminiBaseUrl":"http://127.0.0.1:{p1}",
"anthropicBaseUrl":"http://127.0.0.1:{p2}",
"openAiBaseUrl":"http://127.0.0.1:{p3}/v1",
"liteLlmBaseUrl":"http://127.0.0.1:{p4}",
"awsBedrockEndpoint":"https://vpce.example"}}"#
);
let agents = [agent(&host.root, &state, None)];
for (key, port, proven) in [
("geminiBaseUrl", p1, false),
("anthropicBaseUrl", p2, true),
("openAiBaseUrl", p3, false),
("liteLlmBaseUrl", p4, false),
] {
let mut rec = wired(&host.root, port, proven);
if key == "openAiBaseUrl" {
rec.last_written = Some(format!("http://127.0.0.1:{p3}/v1"));
}
model_relay_endpoints::put_endpoint(&format!("cline:gs:shared:{key}"), rec)
.expect("put");
}
let owner = move |port: u16, _key: &str| match port {
p if p == p3 => PortOwnership::Unreachable,
p if p == p4 => PortOwnership::Foreign,
_ => PortOwnership::Owned,
};
let rows = endpoint_rows_for(&cfg, &agents, &owner);
let state_of = |needle: &str| {
rows.iter()
.find(|r| r.key.ends_with(needle))
.map(|r| r.state.clone())
.unwrap_or_else(|| panic!("no row for {needle}: {rows:#?}"))
};
assert_eq!(state_of("geminiBaseUrl"), EndpointState::NextStart);
assert_eq!(state_of("anthropicBaseUrl"), EndpointState::Active);
assert_eq!(state_of("openAiBaseUrl"), EndpointState::Down);
assert_eq!(state_of("liteLlmBaseUrl"), EndpointState::Foreign);
assert_eq!(state_of("ollamaBaseUrl"), EndpointState::Unwired);
assert_eq!(
state_of("awsBedrockEndpoint"),
EndpointState::Uncovered(UncoveredReason::SignedHost)
);
}
#[test]
fn a_selected_provider_with_no_setting_is_uncovered() {
let host = host();
let agents = [agent(
&host.root,
r#"{"actModeApiProvider":"deepseek"}"#,
Some(r#"{"providers":{}}"#),
)];
let rows = endpoint_rows_for(&config(), &agents, &|_, _| PortOwnership::Owned);
assert_eq!(rows.len(), 1, "{rows:?}");
assert_eq!(rows[0].key, "cline:selected:deepseek");
assert_eq!(
rows[0].state,
EndpointState::Uncovered(UncoveredReason::NoSetting)
);
}
#[test]
fn a_slot_the_editor_may_not_route_by_is_off_until_a_request_proves_it() {
use crate::hooks::model_relay_endpoints::Proof;
let host = host();
let cfg = config();
let block = crate::model_relay::endpoints::endpoint_port_block(cfg.model_relay.port);
let (p1, p2) = (*block.start(), block.start() + 1);
let state = format!(
r#"{{"actModeApiProvider":"anthropic","lastManagedOrganizationId":"org_1",
"anthropicBaseUrl":"http://127.0.0.1:{p1}"}}"#
);
let providers = format!(
r#"{{"providers":{{"deepseek":{{"settings":{{"provider":"deepseek",
"baseUrl":"http://127.0.0.1:{p2}/v1"}}}}}}}}"#
);
let agents = [agent(&host.root, &state, Some(&providers))];
let mut anthropic = wired(&host.root, p1, true);
anthropic.proven_by = Some(Proof::EditorSave);
model_relay_endpoints::put_endpoint("cline:gs:shared:anthropicBaseUrl", anthropic.clone())
.expect("put");
let mut deepseek = wired(&host.root, p2, true);
deepseek.last_written = Some(format!("http://127.0.0.1:{p2}/v1"));
deepseek.file = host
.root
.join("data")
.join("settings")
.join("providers.json");
deepseek.proven_by = Some(Proof::EditorSave);
model_relay_endpoints::put_endpoint("cline:pj:deepseek", deepseek.clone()).expect("put");
let state_of = |needle: &str| {
endpoint_rows_for(&cfg, &agents, &|_, _| PortOwnership::Owned)
.into_iter()
.find(|r| r.key.ends_with(needle))
.map(|r| r.state)
.unwrap_or_else(|| panic!("no row for {needle}"))
};
assert_eq!(
state_of("anthropicBaseUrl"),
EndpointState::Uncovered(UncoveredReason::ManagedOverride)
);
assert_eq!(
state_of("pj:deepseek"),
EndpointState::Uncovered(UncoveredReason::NextBundleOnly)
);
for (key, mut rec) in [
("cline:gs:shared:anthropicBaseUrl", anthropic),
("cline:pj:deepseek", deepseek),
] {
rec.proven_by = Some(Proof::Traffic);
model_relay_endpoints::put_endpoint(key, rec).expect("put");
}
assert_eq!(state_of("anthropicBaseUrl"), EndpointState::Active);
assert_eq!(state_of("pj:deepseek"), EndpointState::Active);
}
#[test]
fn a_misconfigured_slot_is_reported_from_its_record() {
use crate::hooks::model_relay_endpoints::Proof;
let host = host();
let cfg = config();
let port =
*crate::model_relay::endpoints::endpoint_port_block(cfg.model_relay.port).start();
let state = format!(
r#"{{"actModeApiProvider":"gemini","geminiBaseUrl":"http://127.0.0.1:{port}"}}"#
);
let agents = [agent(&host.root, &state, None)];
let mut rec = wired(&host.root, port, true);
rec.proven_by = Some(Proof::EditorSave);
rec.misconfigured_event = Some("evt".into());
model_relay_endpoints::put_endpoint("cline:gs:shared:geminiBaseUrl", rec).expect("put");
let rows = endpoint_rows_for(&cfg, &agents, &|_, _| PortOwnership::Owned);
assert!(
matches!(
&rows[..],
[row] if matches!(row.state, EndpointState::Verdict { code, .. }
if code == crate::error::ERR_MODEL_RELAY_MISCONFIGURED)
),
"{rows:#?}"
);
}
#[test]
fn model_relay_status_reports_each_slot_as_doctor_does() {
use crate::cli::commands::model_relay::{endpoint_severity, endpoint_status, EndpointRow};
use crate::hooks::model_relay_endpoints::Proof;
let host = host();
let mut rec = wired(&host.root, 7601, true);
rec.origin = "https://user:secret@gw.corp.example/".into();
rec.proven_by = Some(Proof::Traffic);
model_relay_endpoints::put_endpoint("cline:gs:shared:geminiBaseUrl", rec).expect("put");
let row = |key: &str, port, state| EndpointRow {
agent: "cline",
key: key.into(),
label: format!("{key} (label)"),
file: None,
port,
state,
};
let rows = vec![
row(
"cline:gs:shared:geminiBaseUrl",
Some(7601),
EndpointState::Active,
),
row(
"cline:gs:shared:ollamaBaseUrl",
Some(7602),
EndpointState::NextStart,
),
row(
"cline:gs:shared:openAiBaseUrl",
Some(7603),
EndpointState::Down,
),
];
let status = serde_json::json!({
"endpoints": [{"key": "cline:gs:shared:geminiBaseUrl", "requests": 3, "contested": true}]
});
let (entries, checks) = endpoint_status(&rows, Some(&status));
let mut report = crate::cli::report::Report::new();
crate::cli::commands::doctor::check_provider_endpoints(&rows, &mut report);
for (entry, check) in entries.iter().zip(report.checks()) {
assert_eq!(entry["state"], check.state.key());
assert_eq!(entry["code"], serde_json::json!(check.code));
assert_eq!(entry["remedy"], serde_json::json!(check.remedy));
}
assert_eq!(entries[0]["state"], "ok");
assert_eq!(entries[1]["state"], "pending");
assert_eq!(entries[1]["code"], crate::error::ERR_MODEL_RELAY_NEXT_START);
let origin = entries[0]["origin"].as_str().expect("origin");
assert!(!origin.contains("secret"), "{origin}");
assert_eq!(entries[0]["proven_by"], "traffic");
assert_eq!(entries[0]["requests"], 3);
assert_eq!(entries[0]["contested"], true);
assert_eq!(entries[1]["contested"], false);
assert_eq!(endpoint_severity(&checks), 2, "a slot nothing serves fails");
assert_eq!(endpoint_severity(&checks[..2]), 1, "a pending slot warns");
assert_eq!(endpoint_severity(&checks[..1]), 0);
}
#[test]
fn only_a_request_since_the_last_write_makes_a_slot_active() {
let host = host();
let cfg = config();
let port =
*crate::model_relay::endpoints::endpoint_port_block(cfg.model_relay.port).start();
let state = format!(
r#"{{"actModeApiProvider":"gemini","geminiBaseUrl":"http://127.0.0.1:{port}"}}"#
);
let gs = host.root.join("data").join("globalState.json");
std::fs::create_dir_all(gs.parent().expect("parent")).expect("mkdir");
std::fs::write(&gs, &state).expect("write");
let endpoints =
ClineProviderEndpoints::at(state_lanes_from(Some(gs.clone()), Some(gs)), None);
let mut rec = wired(&host.root, port, false);
rec.changed_at = 5_000;
model_relay_endpoints::put_endpoint("cline:gs:shared:geminiBaseUrl", rec).expect("put");
let state_with = |last_request: u64| {
let status = serde_json::json!({"endpoints": [{
"key": "cline:gs:shared:geminiBaseUrl",
"requests": 7,
"last_request_unix": last_request,
}]});
super::rows_for_agent(&cfg, &endpoints, Some(&status), &|_, _| {
PortOwnership::Owned
})
.into_iter()
.map(|r| r.state)
.collect::<Vec<_>>()
};
assert_eq!(state_with(4_999), vec![EndpointState::NextStart]);
assert_eq!(state_with(5_000), vec![EndpointState::Active]);
}
#[test]
fn an_isolated_instance_reports_no_slots() {
let host = host();
let mut cfg = config();
cfg.model_relay.own_agent_wiring = Some(false);
let agents = [agent(
&host.root,
r#"{"actModeApiProvider":"ollama"}"#,
None,
)];
assert!(endpoint_rows_for(&cfg, &agents, &|_, _| PortOwnership::Owned).is_empty());
}
}
#[cfg(test)]
mod tests {
use super::{verify_port_ownership, PortOwnership};
use crate::core::egress::test_support::dead_port;
#[test]
fn verify_port_ownership_refuses_closed_and_foreign_ports() {
let closed = dead_port();
assert_eq!(verify_port_ownership(closed), PortOwnership::Unreachable);
let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap();
let foreign = listener.local_addr().unwrap().port();
std::thread::spawn(move || {
use std::io::{Read, Write};
for mut s in listener.incoming().flatten() {
let mut buf = [0u8; 1024];
let _ = s.read(&mut buf);
let body = br#"{"foo":"bar"}"#;
let head = format!(
"HTTP/1.1 200 OK\r\nContent-Type: application/json\r\n\
Content-Length: {}\r\nConnection: close\r\n\r\n",
body.len()
);
let _ = s.write_all(head.as_bytes());
let _ = s.write_all(body);
let _ = s.flush();
}
});
std::thread::sleep(std::time::Duration::from_millis(50));
assert_eq!(verify_port_ownership(foreign), PortOwnership::Foreign);
}
#[tokio::test(flavor = "multi_thread")]
async fn verify_endpoint_ownership_rejects_main_and_sibling_listeners() {
use super::verify_endpoint_ownership;
use crate::model_relay::endpoints::{
normalize_origin, EndpointListeners, EndpointSpec, RelayPorts,
};
use crate::model_relay::{serve_ephemeral, ModelRelayState};
use std::sync::Arc;
let nowhere = reqwest::Url::parse("http://127.0.0.1:9").expect("url");
let factory_base = nowhere.clone();
let listeners = EndpointListeners::new(Arc::new(move |_| {
ModelRelayState::new(factory_base.clone(), 0, 1, &[])
}));
let free = || {
std::net::TcpListener::bind("127.0.0.1:0")
.and_then(|l| l.local_addr())
.expect("free port")
.port()
};
let ports = RelayPorts {
daemon: 1,
main: u16::MAX,
};
let mut slot_ports = Vec::new();
for key in ["test:a", "test:b"] {
let port = free();
listeners
.ensure(EndpointSpec {
key: key.to_string(),
agent: "cline",
family: None,
port,
origin: normalize_origin(&nowhere, &ports).expect("origin"),
})
.await
.expect("endpoint binds");
slot_ports.push(port);
}
let main = serve_ephemeral(Arc::new(ModelRelayState::new(nowhere, 0, 1, &[]))).await;
let dead = dead_port();
let (a, b) = (slot_ports[0], slot_ports[1]);
tokio::task::spawn_blocking(move || {
assert_eq!(verify_endpoint_ownership(a, "test:a"), PortOwnership::Owned);
assert_eq!(
verify_endpoint_ownership(b, "test:a"),
PortOwnership::Foreign,
"a sibling slot's listener is not this slot's"
);
assert_eq!(
verify_endpoint_ownership(main, "test:a"),
PortOwnership::Foreign,
"the main listener is not this slot's"
);
assert_eq!(
verify_endpoint_ownership(dead, "test:a"),
PortOwnership::Unreachable
);
assert_eq!(verify_port_ownership(a), PortOwnership::Owned);
})
.await
.expect("blocking checks");
listeners.shutdown_all().await;
}
#[test]
fn disabled_in_config_short_circuits_before_any_probe() {
use crate::cli::commands::model_relay::{classify_model_relay, ModelRelayState};
let mut cfg = crate::config::Config::defaults();
cfg.model_relay.enabled = false;
assert_eq!(
classify_model_relay(&cfg, "claude-code", None),
ModelRelayState::Disabled
);
assert_eq!(
classify_model_relay(&cfg, "claude-code", Some("http://127.0.0.1:7600")),
ModelRelayState::Disabled
);
assert_eq!(ModelRelayState::Disabled.label(), "disabled");
}
#[test]
fn isolated_instance_short_circuits_before_any_probe() {
use crate::cli::commands::model_relay::{classify_model_relay, ModelRelayState};
let mut cfg = crate::config::Config::defaults();
cfg.model_relay.enabled = true;
cfg.model_relay.port = crate::model_relay::default_model_relay_port() + 1;
assert_eq!(
classify_model_relay(&cfg, "claude-code", None),
ModelRelayState::Isolated
);
}
#[test]
fn status_json_carries_per_agent_wireformat() {
use super::{agent_row_json, AgentRelay, ModelRelayState};
use crate::model_relay::wire_format::WireFormat;
let row = AgentRelay {
agent: "cline",
display_name: "Cline",
wired: Some("http://127.0.0.1:7600/v1".to_string()),
state: ModelRelayState::Wired,
wire_format: Some(WireFormat::OpenAiChatCompletions),
route_hint: "providers.openai-compatible.settings.baseUrl".to_string(),
};
let json = agent_row_json(&row);
assert_eq!(json["agent"], "cline");
assert_eq!(
json["wireformat"], "openai-chat-completions",
"the field must carry the resolved format's wire name"
);
assert_ne!(
json["wireformat"], "unknown",
"`unknown` on a wired plane is a misroute, not a shrug"
);
assert!(
json.get("up").is_none(),
"`up` is a TOP-LEVEL sibling of `agents`, never a member"
);
assert_eq!(
json["wired_to"], "http://127.0.0.1:7600/v1",
"and the shipped fields are untouched"
);
let google = AgentRelay {
wire_format: Some(WireFormat::GoogleGenerateContent),
..row
};
assert_eq!(
agent_row_json(&google)["wireformat"],
"google-generate-content"
);
let editor = AgentRelay {
wire_format: None,
..google
};
assert!(agent_row_json(&editor)["wireformat"].is_null());
}
}