use std::time::Duration;
use super::client::{IntelError, Provider, Transport, resolve};
use super::health::{BreakerConfig, HealthRecord};
const TOKEN_ENV: &str = "AGENTD_INTELLIGENCE_TOKEN";
const TOKEN_ENV_NEUTRAL: &str = "AGENT_INTELLIGENCE_TOKEN";
pub struct Endpoint {
pub(super) transport: Transport,
pub(super) http_path: String,
pub(super) host_header: String,
pub(super) token: Option<String>,
pub(super) provider: Provider,
pub(super) scheme: &'static str,
pub(super) addr: String,
pub(super) extra_headers: Vec<(String, String)>,
pub(super) signer: Option<std::sync::Arc<dyn ::mcp::http::RequestSigner>>,
pub health: HealthRecord,
}
pub struct EndpointList {
eps: Vec<Endpoint>,
active: usize,
breaker: BreakerConfig,
}
fn env(name: &str) -> Option<String> {
std::env::var(name).ok()
}
impl EndpointList {
pub fn parse(uri: &str, default_token: Option<String>) -> Result<EndpointList, IntelError> {
Self::parse_with_env(uri, default_token, &env)
}
pub fn parse_with_env(
uri: &str,
default_token: Option<String>,
env: &dyn Fn(&str) -> Option<String>,
) -> Result<EndpointList, IntelError> {
let provider = Provider::OpenAiCompatible;
let parts: Vec<&str> = uri
.split(',')
.map(str::trim)
.filter(|s| !s.is_empty())
.collect();
if parts.is_empty() {
return Err(IntelError::Unsupported(
"empty intelligence endpoint list".into(),
));
}
let mut eps = Vec::with_capacity(parts.len());
for (i, part) in parts.iter().enumerate() {
let (transport, http_path, host_header) = resolve(part, provider)?;
let token = resolve_token(i, default_token.as_deref(), env)?;
let (scheme, addr) = scheme_and_addr(part);
eps.push(Endpoint {
transport,
http_path,
host_header,
token,
provider,
scheme,
addr,
extra_headers: Vec::new(),
signer: None,
health: HealthRecord::new(),
});
}
Ok(EndpointList {
eps,
active: 0,
breaker: BreakerConfig::default(),
})
}
pub fn set_extra_headers(&mut self, headers: Vec<(String, String)>) {
for e in &mut self.eps {
e.extra_headers = headers.clone();
}
}
pub fn set_signer(&mut self, signer: std::sync::Arc<dyn ::mcp::http::RequestSigner>) {
for e in &mut self.eps {
e.signer = Some(signer.clone());
}
}
pub fn set_provider(&mut self, provider: Provider) {
let default = provider.default_path();
for e in &mut self.eps {
e.provider = provider;
if e.http_path == super::openai::DEFAULT_PATH {
e.http_path = default.to_string();
}
}
}
pub fn len(&self) -> usize {
self.eps.len()
}
pub fn is_empty(&self) -> bool {
self.eps.is_empty()
}
pub fn active(&self) -> usize {
self.active
}
pub fn breaker_config(&self) -> &BreakerConfig {
&self.breaker
}
pub fn ep(&self, idx: usize) -> &Endpoint {
&self.eps[idx]
}
pub fn iter(&self) -> impl Iterator<Item = &Endpoint> {
self.eps.iter()
}
pub fn attempt_order(&self) -> Vec<usize> {
let mut order = Vec::with_capacity(self.eps.len());
if self.eps[self.active].health.available(&self.breaker) {
order.push(self.active);
}
for idx in 0..self.eps.len() {
if idx == self.active {
continue;
}
if self.eps[idx].health.available(&self.breaker) {
order.push(idx);
}
}
order
}
pub fn prefer_lowest_healthy(&mut self) -> Option<usize> {
let target = (0..self.eps.len()).find(|&i| self.eps[i].health.is_up());
if let Some(t) = target
&& t != self.active
{
self.active = t;
return Some(t);
}
None
}
pub fn set_active(&mut self, idx: usize) -> Option<usize> {
if idx != self.active {
self.active = idx;
Some(idx)
} else {
None
}
}
pub fn all_down(&self) -> bool {
self.attempt_order().is_empty()
}
pub fn active_identity(&self) -> (usize, &'static str) {
(self.active, self.eps[self.active].scheme)
}
pub fn body(&self, model: Option<&str>) -> serde_json::Value {
use serde_json::json;
let cfg = &self.breaker;
let endpoints: Vec<serde_json::Value> = self
.eps
.iter()
.enumerate()
.map(|(i, ep)| {
let h = &ep.health;
let mut e = json!({
"index": i,
"transport": ep.scheme,
"addr": ep.addr,
"state": h.state().as_str(),
"active": i == self.active,
"ewma_latency_ms": h.ewma_latency_ms(),
"error_rate": h.error_rate(),
"consec_fail": h.consec_fail(),
});
if let serde_json::Value::Object(m) = &mut e {
if let Some(ms) = h.last_ok_ms_ago() {
m.insert("last_ok_ms_ago".into(), json!(ms));
}
if h.state() == super::health::BreakerState::Open {
if let Some(ms) = h.opened_ms_ago() {
m.insert("opened_ms_ago".into(), json!(ms));
}
m.insert(
"cooldown_ms".into(),
json!(h.cooldown(cfg).as_millis() as u64),
);
m.insert("last_err".into(), json!(h.last_err_kind().as_str()));
}
}
e
})
.collect();
json!({
"active": self.active,
"all_down": self.all_down(),
"model": model,
"endpoints": endpoints,
})
}
}
fn resolve_token(
idx: usize,
default_token: Option<&str>,
env: &dyn Fn(&str) -> Option<String>,
) -> Result<Option<String>, IntelError> {
let (inline_var, file_var, inline_var_n, file_var_n) = if idx == 0 {
(
TOKEN_ENV.to_string(),
format!("{TOKEN_ENV}_FILE"),
TOKEN_ENV_NEUTRAL.to_string(),
format!("{TOKEN_ENV_NEUTRAL}_FILE"),
)
} else {
let n = idx + 1;
(
format!("{TOKEN_ENV}_{n}"),
format!("{TOKEN_ENV}_{n}_FILE"),
format!("{TOKEN_ENV_NEUTRAL}_{n}"),
format!("{TOKEN_ENV_NEUTRAL}_{n}_FILE"),
)
};
if let Some(v) = env(&inline_var_n).or_else(|| env(&inline_var)) {
return Ok(Some(v));
}
if let Some(path) = env(&file_var_n).or_else(|| env(&file_var)) {
let tok = crate::sec::secret::read_token_file(&path).map_err(IntelError::Unsupported)?;
return Ok(Some(tok));
}
if idx == 0 {
return Ok(default_token.map(str::to_string));
}
Ok(None)
}
fn scheme_and_addr(uri: &str) -> (&'static str, String) {
if let Some(rest) = uri.strip_prefix("https://") {
("https", host_only(rest))
} else if let Some(rest) = uri.strip_prefix("http://") {
("http", host_only(rest))
} else {
("unknown", String::new())
}
}
fn host_only(rest: &str) -> String {
rest.split('/').next().unwrap_or(rest).to_string()
}
fn wire_tool_name(name: &str) -> String {
name.chars()
.map(|c| {
if c.is_ascii_alphanumeric() || c == '_' || c == '-' {
c
} else {
'_'
}
})
.collect()
}
impl Endpoint {
#[cfg(feature = "aauth")]
fn aauth_headers(&self, method: &str, path: &str, body: &[u8]) -> Vec<(String, String)> {
match crate::aauth::signer() {
Some(signer) => signer.sign(method, &self.host_header, path, body),
None => Vec::new(),
}
}
#[cfg(not(feature = "aauth"))]
fn aauth_headers(&self, _method: &str, _path: &str, _body: &[u8]) -> Vec<(String, String)> {
Vec::new()
}
pub(super) fn complete_once(
&self,
req: &crate::wire::intel::Request,
timeout: Duration,
trace_id: Option<&str>,
) -> Result<(crate::wire::intel::Response, Duration), IntelError> {
use super::{anthropic, openai};
use crate::net::http;
use std::collections::HashMap;
use std::time::Instant;
use crate::wire::intel::Message;
let dirty = |n: &str| wire_tool_name(n) != n;
let must_sanitize = req.tools.iter().any(|t| dirty(&t.name))
|| req.messages.iter().any(|m| {
matches!(m, Message::Assistant { tool_calls, .. }
if tool_calls.iter().any(|tc| dirty(&tc.name)))
});
let mut wire_to_orig: HashMap<String, String> = HashMap::new();
let owned_req;
let req: &crate::wire::intel::Request = if must_sanitize {
let mut r = req.clone();
for t in &mut r.tools {
let w = wire_tool_name(&t.name);
if w != t.name {
wire_to_orig.insert(w.clone(), t.name.clone());
t.name = w;
}
}
for m in &mut r.messages {
if let Message::Assistant { tool_calls, .. } = m {
for tc in tool_calls {
let w = wire_tool_name(&tc.name);
if w != tc.name {
wire_to_orig.insert(w.clone(), tc.name.clone());
tc.name = w;
}
}
}
}
owned_req = r;
&owned_req
} else {
req
};
use super::bedrock;
let (body, mut headers) = match self.provider {
Provider::OpenAiCompatible => openai::build_request(req, self.token.as_deref()),
Provider::Anthropic => anthropic::build_request(req, self.token.as_deref()),
Provider::Bedrock => bedrock::build_request(req, self.token.as_deref()),
};
let path = self.provider.request_path(&self.http_path, req);
if let Some(tid) = trace_id {
headers.push((
"traceparent".into(),
crate::obs::trace::outbound_traceparent(tid),
));
}
for (k, v) in &self.extra_headers {
headers.push((k.clone(), v.clone()));
}
for (k, v) in self.aauth_headers("POST", &path, &body) {
headers.push((k, v));
}
if let Some(signer) = &self.signer {
for (k, v) in signer.sign("POST", &self.host_header, &path, &body) {
headers.push((k, v));
}
}
let header_refs: Vec<(&str, &str)> = headers
.iter()
.map(|(k, v)| (k.as_str(), v.as_str()))
.collect();
const TRANSIENT_RETRIES: u32 = 2;
let mut attempt: u32 = 0;
let (resp, latency) = loop {
let start = Instant::now();
let mut stream = self.transport.connect(timeout)?;
let resp = http::send(
stream.as_mut(),
&self.host_header,
"POST",
&path,
&header_refs,
&body,
)?;
let latency = start.elapsed();
if resp.is_success() {
break (resp, latency);
}
if super::failover::is_transient_status(resp.status) && attempt < TRANSIENT_RETRIES {
attempt += 1;
std::thread::sleep(Duration::from_millis(250 * (1u64 << (attempt - 1))));
continue;
}
let snippet: String = resp.body_str().chars().take(512).collect();
return Err(IntelError::Http(resp.status, snippet));
};
let mut parsed = match self.provider {
Provider::OpenAiCompatible => openai::parse_response(&resp.body),
Provider::Anthropic => anthropic::parse_response(&resp.body),
Provider::Bedrock => bedrock::parse_response(&resp.body),
}
.map_err(IntelError::Parse)?;
if !wire_to_orig.is_empty() {
for tc in &mut parsed.tool_calls {
if let Some(orig) = wire_to_orig.get(&tc.name) {
tc.name = orig.clone();
}
}
}
Ok((parsed, latency))
}
pub(super) fn discover_models(&self, timeout: Duration) -> Vec<String> {
use super::openai;
use crate::net::http;
if self.provider != Provider::OpenAiCompatible {
return Vec::new();
}
let path = openai::models_path(&self.http_path);
let mut headers: Vec<(String, String)> = Vec::new();
if let Some(tok) = self.token.as_deref() {
headers.push(("authorization".into(), format!("Bearer {tok}")));
}
for (k, v) in self.aauth_headers("GET", &path, &[]) {
headers.push((k, v));
}
let header_refs: Vec<(&str, &str)> = headers
.iter()
.map(|(k, v)| (k.as_str(), v.as_str()))
.collect();
let Ok(mut stream) = self.transport.connect(timeout) else {
return Vec::new();
};
let Ok(resp) = http::send(
stream.as_mut(),
&self.host_header,
"GET",
&path,
&header_refs,
&[],
) else {
return Vec::new();
};
if !resp.is_success() {
return Vec::new();
}
openai::parse_models(&resp.body)
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn extra_headers_apply_to_every_endpoint() {
let mut list = EndpointList::parse(
"https://a.example/v1,https://b.example/v1",
Some("tok".into()),
)
.unwrap();
assert!(list.eps.iter().all(|e| e.extra_headers.is_empty()));
list.set_extra_headers(vec![("X-Team".into(), "ops".into())]);
for e in &list.eps {
assert_eq!(
e.extra_headers,
vec![("X-Team".to_string(), "ops".to_string())]
);
}
}
#[test]
fn wire_tool_name_maps_only_illegal_chars() {
assert_eq!(wire_tool_name("resource.read"), "resource_read");
assert_eq!(wire_tool_name("subagent.spawn"), "subagent_spawn");
assert_eq!(wire_tool_name("math.factorial"), "math_factorial");
assert_eq!(wire_tool_name("get_weather"), "get_weather");
assert_eq!(wire_tool_name("list-files"), "list-files");
assert_eq!(
wire_tool_name("calculate_triangle_area"),
"calculate_triangle_area"
);
assert_eq!(wire_tool_name("a b/c"), "a_b_c");
}
fn env_of<'a>(pairs: &'a [(&'a str, &'a str)]) -> impl Fn(&str) -> Option<String> + 'a {
move |k: &str| {
pairs
.iter()
.find(|(n, _)| *n == k)
.map(|(_, v)| (*v).to_string())
}
}
#[test]
fn comma_list_parses_to_n_endpoints_in_order() {
let env = env_of(&[]);
let list = EndpointList::parse_with_env(
"https://gw-a.example:8443,https://gw-b.example:8444,https://intel.example",
None,
&env,
)
.unwrap();
assert_eq!(list.len(), 3);
assert_eq!(list.ep(0).scheme, "https");
assert_eq!(list.ep(0).addr, "gw-a.example:8443");
assert_eq!(list.ep(1).addr, "gw-b.example:8444");
assert_eq!(list.ep(2).scheme, "https");
assert_eq!(list.active(), 0);
}
#[test]
fn whitespace_around_elements_is_trimmed() {
let env = env_of(&[]);
let list =
EndpointList::parse_with_env(" https://a.example , https://b.example ", None, &env)
.unwrap();
assert_eq!(list.len(), 2);
assert_eq!(list.ep(0).addr, "a.example");
assert_eq!(list.ep(1).addr, "b.example");
}
#[test]
fn empty_list_is_an_error() {
let env = env_of(&[]);
assert!(EndpointList::parse_with_env("", None, &env).is_err());
assert!(EndpointList::parse_with_env(" , ,", None, &env).is_err());
}
#[test]
fn bad_element_scheme_is_an_error() {
let env = env_of(&[]);
let r = EndpointList::parse_with_env("https://a.example,ftp://nope", None, &env);
assert!(matches!(r, Err(IntelError::Unsupported(_))));
for uri in ["unix:/a", "vsock:3:8080", "http://not-loopback.example"] {
let r = EndpointList::parse_with_env(uri, None, &env);
assert!(matches!(r, Err(IntelError::Unsupported(_))), "{uri}");
}
}
#[test]
fn per_endpoint_token_env_resolves_by_position() {
let env = env_of(&[
("AGENTD_INTELLIGENCE_TOKEN", "tok-a"),
("AGENTD_INTELLIGENCE_TOKEN_2", "tok-b"),
]);
let list = EndpointList::parse_with_env("https://a.example,https://b.example", None, &env)
.unwrap();
assert_eq!(list.ep(0).token.as_deref(), Some("tok-a"));
assert_eq!(list.ep(1).token.as_deref(), Some("tok-b"));
}
#[test]
fn endpoint_0_falls_back_to_default_token_when_env_unset() {
let env = env_of(&[]);
let list = EndpointList::parse_with_env(
"https://a.example,https://b.example",
Some("default".into()),
&env,
)
.unwrap();
assert_eq!(list.ep(0).token.as_deref(), Some("default"));
assert_eq!(list.ep(1).token, None);
}
#[test]
fn per_endpoint_env_override_wins_over_default() {
let env = env_of(&[("AGENTD_INTELLIGENCE_TOKEN", "from-env")]);
let list = EndpointList::parse_with_env("https://a.example", Some("default".into()), &env)
.unwrap();
assert_eq!(list.ep(0).token.as_deref(), Some("from-env"));
}
#[test]
fn neutral_token_env_is_accepted_as_an_alias() {
let env = env_of(&[
("AGENT_INTELLIGENCE_TOKEN", "neutral-a"),
("AGENT_INTELLIGENCE_TOKEN_2", "neutral-b"),
]);
let list = EndpointList::parse_with_env("https://a.example,https://b.example", None, &env)
.unwrap();
assert_eq!(list.ep(0).token.as_deref(), Some("neutral-a"));
assert_eq!(list.ep(1).token.as_deref(), Some("neutral-b"));
}
#[test]
fn branded_token_env_wins_over_neutral_on_conflict() {
let env = env_of(&[
("AGENT_INTELLIGENCE_TOKEN", "neutral"),
("AGENTD_INTELLIGENCE_TOKEN", "branded"),
]);
let list = EndpointList::parse_with_env("https://a.example", None, &env).unwrap();
assert_eq!(list.ep(0).token.as_deref(), Some("neutral"));
let env = env_of(&[("AGENTD_INTELLIGENCE_TOKEN", "branded")]);
let list = EndpointList::parse_with_env("https://a.example", None, &env).unwrap();
assert_eq!(list.ep(0).token.as_deref(), Some("branded"));
}
#[test]
fn token_file_variant_reads_from_disk() {
use std::io::Write;
let mut f = tempfile::NamedTempFile::new().unwrap();
writeln!(f, "file-secret").unwrap();
let path = f.path().to_str().unwrap().to_string();
let pairs = [("AGENTD_INTELLIGENCE_TOKEN_2_FILE", path.as_str())];
let env = env_of(&pairs);
let list = EndpointList::parse_with_env("https://a.example,https://b.example", None, &env)
.unwrap();
assert_eq!(list.ep(1).token.as_deref(), Some("file-secret"));
}
#[test]
fn single_element_list_has_inert_failover() {
let env = env_of(&[]);
let list = EndpointList::parse_with_env("https://intel.example", None, &env).unwrap();
assert_eq!(list.len(), 1);
assert_eq!(list.attempt_order(), vec![0]);
assert!(!list.all_down());
}
#[test]
fn attempt_order_skips_open_endpoint_and_snaps_back() {
use super::super::health::ErrKind;
let env = env_of(&[]);
let mut list =
EndpointList::parse_with_env("https://a.example,https://b.example", None, &env)
.unwrap();
let cfg = *list.breaker_config();
for _ in 0..3 {
list.ep(0).health.record_failure(ErrKind::Refused, &cfg);
}
assert_eq!(list.attempt_order(), vec![1]);
assert_eq!(list.prefer_lowest_healthy(), Some(1));
assert_eq!(list.active(), 1);
list.ep(0).health.record_success(Duration::from_millis(5));
assert_eq!(list.prefer_lowest_healthy(), Some(0));
assert_eq!(list.active(), 0);
}
#[test]
fn resource_body_has_health_and_no_url_or_token() {
use super::super::health::ErrKind;
let env = env_of(&[("AGENTD_INTELLIGENCE_TOKEN", "super-secret-tok")]);
let list = EndpointList::parse_with_env(
"https://gw-a.example:8443,https://gw-b.example/v1/secret-path",
None,
&env,
)
.unwrap();
list.ep(0).health.record_success(Duration::from_millis(41));
let cfg = *list.breaker_config();
for _ in 0..3 {
list.ep(1).health.record_failure(ErrKind::Refused, &cfg);
}
let body = list.body(Some("claude-opus-4"));
let text = body.to_string();
assert_eq!(body["active"], 0);
assert_eq!(body["model"], "claude-opus-4");
assert_eq!(body["endpoints"][0]["transport"], "https");
assert_eq!(body["endpoints"][0]["addr"], "gw-a.example:8443");
assert_eq!(body["endpoints"][0]["state"], "closed");
assert_eq!(body["endpoints"][0]["active"], true);
assert_eq!(body["endpoints"][0]["ewma_latency_ms"], 41);
assert_eq!(body["endpoints"][1]["state"], "open");
assert_eq!(body["endpoints"][1]["last_err"], "refused");
assert!(!text.contains("super-secret-tok"), "token leaked: {text}");
assert!(!text.contains("https://"), "full URI leaked: {text}");
assert!(!text.contains("secret-path"), "URL path leaked: {text}");
}
#[test]
fn all_down_when_every_breaker_open() {
use super::super::health::ErrKind;
let env = env_of(&[]);
let list = EndpointList::parse_with_env("https://a.example,https://b.example", None, &env)
.unwrap();
let cfg = *list.breaker_config();
for ep in list.iter() {
for _ in 0..3 {
ep.health.record_failure(ErrKind::Refused, &cfg);
}
}
assert!(list.all_down());
assert!(list.attempt_order().is_empty());
}
}