use crate::{
cancellation::AgentCancellation,
http_body::{DEFAULT_BOUNDED_BODY_MAX_BYTES, read_bounded_response_text},
model_catalog::types::{ModelCatalogEntry, ModelCatalogServiceTier},
providers::{
HttpRequest,
transport::{
join_provider_worker_with_timeout, provider_client,
read_bounded_redacted_response_body, sanitize_provider_error_url,
},
},
thinking::ThinkingLevel,
};
use serde_json::Value;
use std::{
sync::{
atomic::{AtomicUsize, Ordering},
mpsc,
},
thread,
time::Duration,
};
const MODEL_CATALOG_WORKER_SHUTDOWN_TIMEOUT: Duration = Duration::from_millis(100);
const MODEL_CATALOG_WORKER_LIMIT: usize = 32;
static MODEL_CATALOG_WORKERS: AtomicUsize = AtomicUsize::new(0);
struct ModelCatalogWorkerSlot<'a> {
workers: &'a AtomicUsize,
}
impl Drop for ModelCatalogWorkerSlot<'_> {
fn drop(&mut self) {
self.workers.fetch_sub(1, Ordering::SeqCst);
}
}
fn reserve_model_catalog_worker_slot(workers: &AtomicUsize) -> Option<ModelCatalogWorkerSlot<'_>> {
let mut current = workers.load(Ordering::SeqCst);
loop {
if current >= MODEL_CATALOG_WORKER_LIMIT {
return None;
}
match workers.compare_exchange(current, current + 1, Ordering::SeqCst, Ordering::SeqCst) {
Ok(_) => return Some(ModelCatalogWorkerSlot { workers }),
Err(next) => current = next,
}
}
}
fn reserve_production_model_catalog_worker_slot() -> Option<ModelCatalogWorkerSlot<'static>> {
reserve_model_catalog_worker_slot(&MODEL_CATALOG_WORKERS)
}
pub(crate) const MODEL_CATALOG_SUCCESS_BODY_MAX_BYTES: u64 = DEFAULT_BOUNDED_BODY_MAX_BYTES;
pub(crate) fn fetch_model_catalog_response_text_cancellable(
request: HttpRequest,
provider_id: &str,
cancellation: &AgentCancellation,
) -> anyhow::Result<String> {
cancellation.check()?;
let provider_id = provider_id.to_string();
let (sender, receiver) = mpsc::sync_channel(1);
let worker_slot = reserve_production_model_catalog_worker_slot().ok_or_else(|| anyhow::anyhow!("model discovery worker limit reached ({MODEL_CATALOG_WORKER_LIMIT}); retry after stalled model discovery requests finish"))?;
let worker = thread::Builder::new()
.name("model-catalog-fetch".to_string())
.spawn(move || {
let _worker_slot = worker_slot;
let _ = sender.send(fetch_model_catalog_response_text_blocking(
request,
&provider_id,
));
})
.map_err(|error| anyhow::anyhow!("model discovery worker spawn failed: {error}"))?;
const CANCEL_POLL_INTERVAL: Duration = Duration::from_millis(25);
loop {
if cancellation.is_canceled() {
drop(receiver);
join_provider_worker_with_timeout(worker, MODEL_CATALOG_WORKER_SHUTDOWN_TIMEOUT);
return match cancellation.check() {
Ok(()) => anyhow::bail!("model discovery canceled"),
Err(error) => Err(error),
};
}
match receiver.recv_timeout(CANCEL_POLL_INTERVAL) {
Ok(result) => {
join_provider_worker_with_timeout(worker, MODEL_CATALOG_WORKER_SHUTDOWN_TIMEOUT);
return result;
}
Err(mpsc::RecvTimeoutError::Timeout) => continue,
Err(mpsc::RecvTimeoutError::Disconnected) => {
join_provider_worker_with_timeout(worker, MODEL_CATALOG_WORKER_SHUTDOWN_TIMEOUT);
cancellation.check()?;
anyhow::bail!("model discovery worker disconnected")
}
}
}
}
fn fetch_model_catalog_response_text_blocking(
request: HttpRequest,
provider_id: &str,
) -> anyhow::Result<String> {
let client = provider_client()?;
let method = reqwest::Method::from_bytes(request.method.as_bytes()).map_err(|error| {
anyhow::anyhow!(
"{provider_id} model discovery invalid HTTP method '{}': {error}",
request.method
)
})?;
let endpoint = sanitize_provider_error_url(&request.url);
let mut builder = client.request(method, &request.url);
for (name, value) in &request.headers {
builder = builder.header(name, value);
}
let response = builder.send().map_err(|error| {
let error = error.without_url();
anyhow::anyhow!("{provider_id} model discovery request failed for {endpoint}: {error}")
})?;
if !response.status().is_success() {
let status = response.status();
let body = read_bounded_redacted_response_body(response);
anyhow::bail!(
"{provider_id} model discovery failed for {endpoint} with status {status}: {body}"
);
}
read_bounded_response_text(response, MODEL_CATALOG_SUCCESS_BODY_MAX_BYTES).map_err(|error| {
let error = error.to_string();
if error.contains("response exceeded") {
anyhow::anyhow!(
"{provider_id} model discovery response exceeded {MODEL_CATALOG_SUCCESS_BODY_MAX_BYTES} bytes"
)
} else {
anyhow::anyhow!("{provider_id} model discovery response read failed: {error}")
}
})
}
pub(crate) fn parse_openai_compatible_model_catalog_response(
provider_id: &str,
text: &str,
) -> anyhow::Result<Vec<ModelCatalogEntry>> {
let value: Value = serde_json::from_str(text)?;
let data = value.get("data").and_then(Value::as_array).ok_or_else(|| {
anyhow::anyhow!("{provider_id} model discovery response missing data array")
})?;
let mut entries = Vec::new();
for item in data {
let Some(model) = item
.get("id")
.and_then(Value::as_str)
.map(str::trim)
.filter(|id| !id.is_empty())
else {
continue;
};
entries.push(ModelCatalogEntry::new(provider_id, model));
}
if entries.is_empty() {
anyhow::bail!("{provider_id} model discovery returned no usable models");
}
Ok(entries)
}
pub(crate) fn parse_codex_model_catalog_response(
text: &str,
) -> anyhow::Result<Vec<ModelCatalogEntry>> {
let value: Value = serde_json::from_str(text)?;
let models = value
.get("models")
.and_then(Value::as_array)
.ok_or_else(|| {
anyhow::anyhow!("openai-codex model discovery response missing models array")
})?;
let mut entries = Vec::new();
for model in models {
let Some(slug) = model.get("slug").and_then(Value::as_str).map(str::trim) else {
continue;
};
if slug.is_empty() {
continue;
}
if model
.get("supported_in_api")
.and_then(Value::as_bool)
.is_some_and(|supported| !supported)
{
continue;
}
if model
.get("visibility")
.and_then(Value::as_str)
.is_some_and(is_hidden_visibility)
{
continue;
}
let mut entry = ModelCatalogEntry::new_codex(slug);
let (reasoning_efforts, supports_reasoning) = parse_codex_reasoning_metadata(model);
entry.reasoning_efforts = reasoning_efforts;
entry.supports_reasoning = supports_reasoning;
entry.service_tiers = parse_codex_service_tiers(model);
entry.display_name = model
.get("display_name")
.and_then(Value::as_str)
.filter(|value| !value.is_empty())
.map(str::to_string);
entry.description = model
.get("description")
.and_then(Value::as_str)
.filter(|value| !value.is_empty())
.map(str::to_string);
entry.context_window = model
.get("context_window")
.or_else(|| model.get("max_context_window"))
.and_then(Value::as_u64);
entry.max_context_window = model.get("max_context_window").and_then(Value::as_u64);
entries.push(entry);
}
if entries.is_empty() {
anyhow::bail!("openai-codex model discovery returned no usable models");
}
Ok(entries)
}
fn parse_codex_reasoning_metadata(model: &Value) -> (Option<Vec<ThinkingLevel>>, Option<bool>) {
let Some(levels) = model
.get("supported_reasoning_levels")
.and_then(Value::as_array)
else {
return (None, None);
};
let recognized = levels
.iter()
.filter_map(|level| level.get("effort").and_then(Value::as_str))
.filter_map(|effort| {
let effort = effort.trim();
(effort.eq_ignore_ascii_case("none"))
.then_some(ThinkingLevel::Default)
.or_else(|| effort.parse().ok())
})
.collect::<Vec<_>>();
let efforts =
(!recognized.is_empty()).then(|| crate::thinking::normalize_thinking_levels(&recognized));
(efforts, Some(!recognized.is_empty()))
}
fn parse_codex_service_tiers(model: &Value) -> Option<Vec<ModelCatalogServiceTier>> {
if let Some(service_tiers) = model.get("service_tiers") {
return Some(
service_tiers
.as_array()
.into_iter()
.flatten()
.filter_map(parse_codex_service_tier)
.collect(),
);
}
model.get("additional_speed_tiers").map(|tiers| {
tiers
.as_array()
.into_iter()
.flatten()
.filter_map(Value::as_str)
.filter(|tier| !tier.trim().is_empty())
.map(|tier| ModelCatalogServiceTier {
id: tier.to_string(),
name: tier.to_string(),
description: None,
})
.collect()
})
}
fn parse_codex_service_tier(value: &Value) -> Option<ModelCatalogServiceTier> {
let object = value.as_object()?;
let id = object.get("id").and_then(Value::as_str)?;
let name = object.get("name").and_then(Value::as_str)?;
if id.trim().is_empty() || name.trim().is_empty() {
return None;
}
let description = match object.get("description") {
None | Some(Value::Null) => None,
Some(Value::String(description)) => Some(description.to_string()),
Some(_) => return None,
};
Some(ModelCatalogServiceTier {
id: id.to_string(),
name: name.to_string(),
description,
})
}
fn is_hidden_visibility(visibility: &str) -> bool {
let visibility = visibility.trim();
visibility.eq_ignore_ascii_case("hide") || visibility.eq_ignore_ascii_case("hidden")
}
pub fn extract_chatgpt_account_id(access_token: &str) -> anyhow::Result<String> {
crate::config::extract_chatgpt_account_id_from_jwt(access_token)
}
#[cfg(test)]
mod tests {
use super::{
MODEL_CATALOG_WORKER_LIMIT, parse_codex_model_catalog_response,
reserve_model_catalog_worker_slot,
};
use serde_json::json;
use std::sync::atomic::AtomicUsize;
#[test]
fn codex_catalog_parses_current_and_legacy_service_tiers_with_explicit_presence() {
let current = parse_codex_model_catalog_response(
&json!({
"models": [{
"slug": "gpt-5.5",
"service_tiers": [
{"id": "priority", "name": "FAST", "description": "low latency"},
{"id": 7, "name": "fast"},
{"id": "", "name": "fast"},
{"id": "ignored", "name": "fast", "description": false}
],
"additional_speed_tiers": ["legacy"]
}]
})
.to_string(),
)
.unwrap();
assert_eq!(
current[0].service_tiers,
Some(vec![crate::model_catalog::types::ModelCatalogServiceTier {
id: "priority".to_string(),
name: "FAST".to_string(),
description: Some("low latency".to_string()),
}])
);
let explicit_empty = parse_codex_model_catalog_response(
&json!({
"models": [{
"slug": "gpt-5.5",
"service_tiers": [],
"additional_speed_tiers": ["fast"]
}]
})
.to_string(),
)
.unwrap();
assert_eq!(explicit_empty[0].service_tiers, Some(Vec::new()));
let legacy = parse_codex_model_catalog_response(
&json!({
"models": [{"slug": "gpt-5.5", "additional_speed_tiers": ["fast", "", 4]}]
})
.to_string(),
)
.unwrap();
assert_eq!(
legacy[0].service_tiers,
Some(vec![crate::model_catalog::types::ModelCatalogServiceTier {
id: "fast".to_string(),
name: "fast".to_string(),
description: None,
}])
);
let absent = parse_codex_model_catalog_response(
&json!({"models": [{"slug": "gpt-5.5"}]}).to_string(),
)
.unwrap();
assert_eq!(absent[0].service_tiers, None);
}
#[test]
fn codex_catalog_parses_provider_supported_reasoning_levels() {
let supported = ["low", "medium", "high", "xhigh", "max", "ultra"]
.into_iter()
.map(|effort| json!({"effort": effort}))
.collect::<Vec<_>>();
let entries = parse_codex_model_catalog_response(
&json!({
"models": [
{"slug": "gpt-6-astra", "supported_reasoning_levels": supported.clone()},
{"slug": "gpt-5.6-sol", "supported_reasoning_levels": supported}
]
})
.to_string(),
)
.unwrap();
assert_eq!(
entries
.iter()
.map(|entry| entry.id.as_str())
.collect::<Vec<_>>(),
vec!["openai-codex/gpt-6-astra", "openai-codex/gpt-5.6-sol"]
);
for entry in entries {
assert_eq!(
entry.reasoning_efforts,
Some(crate::thinking::ThinkingLevel::ALL.to_vec())
);
assert_eq!(entry.supports_reasoning, Some(true));
}
}
#[test]
fn codex_catalog_requires_recognized_reasoning_levels_for_support() {
let cases = [
("gpt-5", json!([{"effort": "ultra"}])),
(
"gpt-malformed-effort",
json!([{"effort": 42}, {"wrong_field": "low"}, null]),
),
("gpt-empty-efforts", json!([])),
];
let models = cases
.iter()
.map(|(slug, levels)| json!({"slug": slug, "supported_reasoning_levels": levels}))
.collect::<Vec<_>>();
let entries =
parse_codex_model_catalog_response(&json!({"models": models}).to_string()).unwrap();
for entry in entries {
assert_eq!(entry.reasoning_efforts, None);
assert_eq!(entry.supports_reasoning, Some(false));
}
}
#[test]
fn model_catalog_worker_slots_are_hard_capped_and_released() {
let workers = AtomicUsize::new(0);
let slots = (0..MODEL_CATALOG_WORKER_LIMIT)
.map(|_| reserve_model_catalog_worker_slot(&workers).expect("catalog worker slot"))
.collect::<Vec<_>>();
assert!(reserve_model_catalog_worker_slot(&workers).is_none());
drop(slots);
assert!(reserve_model_catalog_worker_slot(&workers).is_some());
}
}