magi-code 0.96.2

Repository-aware CLI coding agent for terminal work
Documentation
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);
            // ponytail: blocking catalog HTTP cannot be interrupted; bound-join then detach until reqwest timeout.
            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 {
                // Legacy speed names are not Codex wire-tier identifiers.
                id: if tier.trim().eq_ignore_ascii_case("fast") {
                    "priority".to_string()
                } else {
                    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::auth::extract_chatgpt_account_id_from_jwt(access_token)
}