"""Unified LLM client for GA text generation.

Supports OpenAI, Google Gemini, Anthropic, and Ollama (local) as providers.
Provider and model are set globally via environment variables:

  LLM_PROVIDER      - "openai" (default), "gemini", "anthropic", or "ollama"
  LLM_MODEL         - e.g. "gpt-4o", "gemini-3-flash-preview",
                      "claude-sonnet-4-6", "qwen2.5:7b"
  OLLAMA_BASE_URL   - Ollama server URL (default: http://localhost:11434/v1)

Key pools (comma-separated env vars):
  POOL_OPENAI_KEYS      - rotate across multiple OpenAI keys
  POOL_GEMINI_API_KEYS  - rotate across multiple Gemini keys
  POOL_ANTHROPIC_KEYS   - rotate across multiple Anthropic keys
Falls back to OPENAI_API_KEY / GEMINI_API_KEY / ANTHROPIC_API_KEY when
pools are not set.
"""

import base64
import json
import logging
import os
import re
import threading
from pathlib import Path

import anthropic as _anthropic
import dotenv
import openai as _openai
import requests as _requests

from key_pool import KeyPool

dotenv.load_dotenv()

_log = logging.getLogger(__name__)

# ------------------------------------------------------------------
# Globals
# ------------------------------------------------------------------

_CLIENTS: dict = {}
_CLIENT_LOCK = threading.Lock()

LLM_PROVIDER = os.environ.get("LLM_PROVIDER", "gemini")
LLM_MODEL = os.environ.get("LLM_MODEL", "gemini-3-flash-preview")

_HTTP_TOO_MANY_REQUESTS = 429

# Gemini "thinking" models burn tokens on reasoning before output.
# Override max_output_tokens to ensure enough room for the actual response.
GEMINI_MIN_OUTPUT_TOKENS = 4096

# ------------------------------------------------------------------
# Key pools
# ------------------------------------------------------------------

_openai_pool = KeyPool("POOL_OPENAI_KEYS", "OPENAI_API_KEY", cooldown_seconds=60.0)
_gemini_pool = KeyPool("POOL_GEMINI_API_KEYS", "GEMINI_API_KEY", cooldown_seconds=60.0)
_anthropic_pool = KeyPool(
    "POOL_ANTHROPIC_KEYS", "ANTHROPIC_API_KEY", cooldown_seconds=60.0
)


def get_pool_size() -> int:
    """Return number of keys in the active provider's pool (min 1)."""
    if LLM_PROVIDER == "openai":
        return max(_openai_pool.size, 1)
    if LLM_PROVIDER == "gemini":
        return max(_gemini_pool.size, 1)
    if LLM_PROVIDER == "anthropic":
        return max(_anthropic_pool.size, 1)
    return 1


# ------------------------------------------------------------------
# Rate-limit sentinel
# ------------------------------------------------------------------


class _RateLimitedError(Exception):  # noqa: N818
    """Raised internally when a provider returns 429."""

    def __init__(self, key: str, headers: dict | None = None) -> None:
        """Store the key and response headers for cooldown bookkeeping."""
        self.key = key
        self.headers = headers or {}
        super().__init__(f"Rate limited: ...{key[-4:]}")


# ------------------------------------------------------------------
# Provider helpers (per-key, thread-safe)
# ------------------------------------------------------------------


def _get_openai_client(api_key: str):
    """Return a cached OpenAI client for *api_key*."""
    cache_key = f"openai_{api_key[-8:]}"
    with _CLIENT_LOCK:
        if cache_key not in _CLIENTS:
            _CLIENTS[cache_key] = _openai.OpenAI(api_key=api_key)
        return _CLIENTS[cache_key]


def _get_ollama_client():
    """Return a cached Ollama client (OpenAI-compatible local server)."""
    with _CLIENT_LOCK:
        if "ollama" not in _CLIENTS:
            base_url = os.environ.get("OLLAMA_BASE_URL", "http://localhost:11434/v1")
            _CLIENTS["ollama"] = _openai.OpenAI(base_url=base_url, api_key="ollama")
        return _CLIENTS["ollama"]


def _call_openai(api_key, model, system_message, user_message, temperature, max_tokens):
    """Call OpenAI chat completions. Raises _RateLimitedError on 429."""
    client = _get_openai_client(api_key)
    try:
        resp = client.chat.completions.create(
            model=model,
            temperature=temperature,
            max_tokens=max_tokens,
            messages=[
                {"role": "system", "content": system_message},
                {"role": "user", "content": user_message},
            ],
            response_format={"type": "json_object"},
        )
        return resp.choices[0].message.content
    except _openai.RateLimitError as e:
        headers = {}
        if hasattr(e, "response") and e.response is not None:
            h = e.response.headers
            for name in (
                "retry-after",
                "x-ratelimit-reset-requests",
                "x-ratelimit-reset-tokens",
            ):
                if name in h:
                    headers[name] = h[name]
        raise _RateLimitedError(api_key, headers) from e


def _call_gemini(
    api_key, model, system_message, user_message, temperature, max_output_tokens
):  # noqa: PLR0913
    """Call Gemini REST API directly (thread-safe, no global state).

    Raises _RateLimitedError on 429.
    """
    url = f"https://generativelanguage.googleapis.com/v1beta/models/{model}:generateContent"
    body = {
        "systemInstruction": {"parts": [{"text": system_message}]},
        "contents": [{"parts": [{"text": user_message}]}],
        "generationConfig": {
            "temperature": temperature,
            "maxOutputTokens": max_output_tokens,
            "responseMimeType": "application/json",
        },
    }
    resp = _requests.post(url, params={"key": api_key}, json=body, timeout=120)
    if resp.status_code == _HTTP_TOO_MANY_REQUESTS:
        headers = {}
        for name in (
            "retry-after",
            "x-ratelimit-reset-requests",
            "x-ratelimit-reset-tokens",
        ):
            if name in resp.headers:
                headers[name] = resp.headers[name]
        raise _RateLimitedError(api_key, headers)
    resp.raise_for_status()
    data = resp.json()
    candidates = data.get("candidates", [])
    if not candidates:
        raise ValueError(f"Gemini returned no candidates: {json.dumps(data)[:300]}")
    parts = candidates[0].get("content", {}).get("parts", [])
    if not parts:
        finish = candidates[0].get("finishReason", "UNKNOWN")
        raise ValueError(f"Gemini candidate has no content (finishReason={finish})")
    return parts[0]["text"]


def _call_openai_vision(  # noqa: PLR0913
    api_key, model, system_message, user_message, temperature, max_tokens, image_paths
):
    """Call OpenAI chat completions with vision. Raises _RateLimitedError on 429."""
    client = _get_openai_client(api_key)
    content = []
    for path in image_paths:
        data = base64.b64encode(path.read_bytes()).decode()
        content.append(
            {
                "type": "image_url",
                "image_url": {"url": f"data:image/jpeg;base64,{data}"},
            }
        )
    content.append({"type": "text", "text": user_message})
    try:
        resp = client.chat.completions.create(
            model=model,
            temperature=temperature,
            max_tokens=max_tokens,
            messages=[
                {"role": "system", "content": system_message},
                {"role": "user", "content": content},
            ],
        )
        raw = resp.choices[0].message.content
        return re.sub(r"```(?:json)?|```", "", raw).strip()
    except _openai.RateLimitError as e:
        headers = {}
        if hasattr(e, "response") and e.response is not None:
            h = e.response.headers
            for name in (
                "retry-after",
                "x-ratelimit-reset-requests",
                "x-ratelimit-reset-tokens",
            ):
                if name in h:
                    headers[name] = h[name]
        raise _RateLimitedError(api_key, headers) from e


def _call_gemini_vision(  # noqa: PLR0913
    api_key,
    model,
    system_message,
    user_message,
    temperature,
    max_output_tokens,
    image_paths,
):
    """Call Gemini REST API with vision (inlineData). Raises _RateLimitedError on 429."""
    url = f"https://generativelanguage.googleapis.com/v1beta/models/{model}:generateContent"
    parts = []
    for path in image_paths:
        data = base64.b64encode(path.read_bytes()).decode()
        parts.append({"inlineData": {"mimeType": "image/jpeg", "data": data}})
    parts.append({"text": user_message})
    body = {
        "systemInstruction": {"parts": [{"text": system_message}]},
        "contents": [{"parts": parts}],
        "generationConfig": {
            "temperature": temperature,
            "maxOutputTokens": max_output_tokens,
            "responseMimeType": "application/json",
        },
    }
    resp = _requests.post(url, params={"key": api_key}, json=body, timeout=120)
    if resp.status_code == _HTTP_TOO_MANY_REQUESTS:
        headers = {
            name: resp.headers[name]
            for name in (
                "retry-after",
                "x-ratelimit-reset-requests",
                "x-ratelimit-reset-tokens",
            )
            if name in resp.headers
        }
        raise _RateLimitedError(api_key, headers)
    resp.raise_for_status()
    data = resp.json()
    candidates = data.get("candidates", [])
    if not candidates:
        raise ValueError(f"Gemini returned no candidates: {json.dumps(data)[:300]}")
    content_parts = candidates[0].get("content", {}).get("parts", [])
    if not content_parts:
        finish = candidates[0].get("finishReason", "UNKNOWN")
        raise ValueError(f"Gemini candidate has no content (finishReason={finish})")
    return content_parts[0]["text"]


_AUDIO_MIME = {
    ".m4a": "audio/mp4",
    ".mp4": "audio/mp4",
    ".aac": "audio/aac",
    ".mp3": "audio/mp3",
    ".wav": "audio/wav",
    ".ogg": "audio/ogg",
    ".flac": "audio/flac",
}


def _call_gemini_audio(  # noqa: PLR0913
    api_key,
    model,
    system_message,
    user_message,
    temperature,
    max_output_tokens,
    audio_path,
):
    """Call Gemini REST with an inline audio file; return plain text. 429 -> _RateLimitedError."""
    url = f"https://generativelanguage.googleapis.com/v1beta/models/{model}:generateContent"
    mime = _AUDIO_MIME.get(audio_path.suffix.lower(), "audio/mp4")
    data = base64.b64encode(audio_path.read_bytes()).decode()
    body = {
        "systemInstruction": {"parts": [{"text": system_message}]},
        "contents": [
            {
                "parts": [
                    {"inlineData": {"mimeType": mime, "data": data}},
                    {"text": user_message},
                ]
            }
        ],
        "generationConfig": {
            "temperature": temperature,
            "maxOutputTokens": max_output_tokens,
        },
    }
    resp = _requests.post(url, params={"key": api_key}, json=body, timeout=300)
    if resp.status_code == _HTTP_TOO_MANY_REQUESTS:
        headers = {
            name: resp.headers[name]
            for name in (
                "retry-after",
                "x-ratelimit-reset-requests",
                "x-ratelimit-reset-tokens",
            )
            if name in resp.headers
        }
        raise _RateLimitedError(api_key, headers)
    resp.raise_for_status()
    data = resp.json()
    candidates = data.get("candidates", [])
    if not candidates:
        raise ValueError(f"Gemini returned no candidates: {json.dumps(data)[:300]}")
    content_parts = candidates[0].get("content", {}).get("parts", [])
    if not content_parts:
        finish = candidates[0].get("finishReason", "UNKNOWN")
        raise ValueError(f"Gemini candidate has no content (finishReason={finish})")
    return content_parts[0]["text"]


def transcribe_audio(
    audio_path: Path,
    system_message: str,
    user_message: str,
    temperature: float = 0.0,
    max_output_tokens: int = 8192,
    model: str | None = None,
) -> str | None:
    """Transcribe one audio file via the Gemini key pool. Return text or None.

    Rotates across POOL_GEMINI_API_KEYS with 429 cooldown handling, mirroring
    call_llm_chat_json but returning the raw transcript text (no JSON parsing).
    """
    model = model or LLM_MODEL
    pool = _gemini_pool
    out_tokens = max(max_output_tokens, GEMINI_MIN_OUTPUT_TOKENS)
    max_attempts = pool.size + 3
    max_errors = 2
    errors = 0
    for _attempt in range(max_attempts):
        try:
            key = pool.get_key()
            if not key:
                raise ValueError("No Gemini API keys configured")
            text = _call_gemini_audio(
                key,
                model,
                system_message,
                user_message,
                temperature,
                out_tokens,
                audio_path,
            )
            pool.clear_cooldown(key)
            if text and text.strip():
                return text.strip()
            errors += 1
        except _RateLimitedError as e:
            pool.mark_rate_limited(e.key, headers=e.headers)
            continue
        except (ValueError, OSError, _requests.RequestException) as e:
            errors += 1
            if errors > max_errors:
                _log.warning("audio transcribe failed (%s): %s", audio_path.name, e)
                return None
    return None


def _call_anthropic_vision(  # noqa: PLR0913
    api_key, model, system_message, user_message, temperature, max_tokens, image_paths
):
    """Call Anthropic Messages API with vision. Raises _RateLimitedError on 429."""
    client = _get_anthropic_client(api_key)
    content = []
    for path in image_paths:
        data = base64.b64encode(path.read_bytes()).decode()
        content.append(
            {
                "type": "image",
                "source": {"type": "base64", "media_type": "image/jpeg", "data": data},
            }
        )
    content.append({"type": "text", "text": user_message})
    try:
        msg = client.messages.create(
            model=model,
            max_tokens=max_tokens,
            temperature=temperature,
            system=system_message,
            messages=[{"role": "user", "content": content}],
        )
        raw = msg.content[0].text
        return re.sub(r"```(?:json)?|```", "", raw).strip()
    except _anthropic.RateLimitError as e:
        headers = {}
        if hasattr(e, "response") and e.response is not None:
            h = e.response.headers
            for name in (
                "retry-after",
                "x-ratelimit-reset-requests",
                "x-ratelimit-reset-tokens",
            ):
                if name in h:
                    headers[name] = h[name]
        raise _RateLimitedError(api_key, headers) from e


def _call_ollama(model, system_message, user_message, temperature, max_tokens):
    """Call local Ollama (no key rotation needed)."""
    client = _get_ollama_client()
    resp = client.chat.completions.create(
        model=model,
        temperature=temperature,
        max_tokens=max_tokens,
        messages=[
            {"role": "system", "content": system_message},
            {"role": "user", "content": user_message},
        ],
        response_format={"type": "json_object"},
    )
    return resp.choices[0].message.content


def _get_anthropic_client(api_key: str):
    """Return a cached Anthropic client for *api_key*."""
    cache_key = f"anthropic_{api_key[-8:]}"
    with _CLIENT_LOCK:
        if cache_key not in _CLIENTS:
            _CLIENTS[cache_key] = _anthropic.Anthropic(api_key=api_key)
        return _CLIENTS[cache_key]


def _call_anthropic(api_key, model, system_message, user_message, temperature, max_tokens):  # noqa: PLR0913  # fmt: skip
    """Call Anthropic Messages API. Raises _RateLimitedError on 429."""
    client = _get_anthropic_client(api_key)
    try:
        msg = client.messages.create(
            model=model,
            max_tokens=max_tokens,
            temperature=temperature,
            system=system_message,
            messages=[{"role": "user", "content": user_message}],
        )
        return re.sub(r"```(?:json)?|```", "", msg.content[0].text).strip()
    except _anthropic.RateLimitError as e:
        headers = {}
        if hasattr(e, "response") and e.response is not None:
            h = e.response.headers
            for name in (
                "retry-after",
                "x-ratelimit-reset-requests",
                "x-ratelimit-reset-tokens",
            ):
                if name in h:
                    headers[name] = h[name]
        raise _RateLimitedError(api_key, headers) from e


# ------------------------------------------------------------------
# JSON repair
# ------------------------------------------------------------------


def _repair_truncated_json(text: str) -> dict | None:
    """Attempt to salvage a truncated JSON response (e.g. from max_tokens cutoff).

    Handles the common case: {"variants": ["text1", "text2...
    """
    if not text or not text.strip():
        return None

    match = re.search(r'"variants"\s*:\s*\[', text)
    if not match:
        return None

    after_bracket = text[match.end() :]
    complete_strings = re.findall(r'"((?:[^"\\]|\\.)*)"', after_bracket)

    if complete_strings:
        return {"variants": complete_strings}

    return None


# ------------------------------------------------------------------
# Main entry point
# ------------------------------------------------------------------


def _parse_json_response(
    raw: str, retries: int, errors: int, label: str
) -> tuple[dict | None, int, bool]:
    """Parse the first complete JSON object from *raw*.

    Returns (parsed_obj, updated_errors, should_continue).
    """
    try:
        obj, _ = json.JSONDecoder().raw_decode(raw.strip())
        return obj, errors, False
    except json.JSONDecodeError:
        repaired = _repair_truncated_json(raw)
        if repaired:
            return repaired, errors, False
        errors += 1
        if errors > retries:
            _log.warning("LLM JSON parse failed after %d attempts (%s)", errors, label)
            return None, errors, False
        return None, errors, True  # continue


def call_llm_chat_json(
    prompt_cfg: dict,
    system_message: str,
    user_message: str,
    provider: str | None = None,
    model: str | None = None,
) -> dict | None:  # noqa: C901, PLR0912, PLR0915
    """Call an LLM and return parsed JSON dict, or None on failure.

    Provider and model come from LLM_PROVIDER / LLM_MODEL env vars by default.
    Pass *provider* and *model* to override for this call only.

    Args:
        prompt_cfg: dict with keys temperature (float), max_tokens (int),
            and optional retry_on_invalid_json (int, default 2).
        system_message: System prompt text.
        user_message: User prompt text.
        provider: LLM provider override ("openai", "gemini", "anthropic"). Defaults to LLM_PROVIDER.
        model: Model name override. Defaults to LLM_MODEL.

    Returns:
        Parsed JSON dict, or None on repeated failure.
    """
    provider = provider or LLM_PROVIDER
    model = model or LLM_MODEL
    retries = prompt_cfg.get("retry_on_invalid_json", 2)
    temperature = prompt_cfg["temperature"]
    max_tokens = prompt_cfg["max_tokens"]

    pool = {
        "openai": _openai_pool,
        "gemini": _gemini_pool,
        "anthropic": _anthropic_pool,
    }.get(provider)
    pool_sz = pool.size if pool else 1
    max_attempts = (retries + 1) + pool_sz
    errors = 0

    for _attempt in range(max_attempts):
        try:
            if provider == "openai":
                key = pool.get_key()
                if not key:
                    raise ValueError("No OpenAI API keys configured")
                raw = _call_openai(
                    key, model, system_message, user_message, temperature, max_tokens
                )
                pool.clear_cooldown(key)

            elif provider == "gemini":
                key = pool.get_key()
                if not key:
                    raise ValueError("No Gemini API keys configured")
                output_tokens = max(max_tokens, GEMINI_MIN_OUTPUT_TOKENS)
                raw = _call_gemini(
                    key, model, system_message, user_message, temperature, output_tokens
                )
                pool.clear_cooldown(key)

            elif provider == "anthropic":
                key = pool.get_key()
                if not key:
                    raise ValueError("No Anthropic API keys configured")
                raw = _call_anthropic(
                    key, model, system_message, user_message, temperature, max_tokens
                )
                pool.clear_cooldown(key)

            elif provider == "ollama":
                raw = _call_ollama(
                    model, system_message, user_message, temperature, max_tokens
                )

            else:
                raise ValueError(f"Unknown LLM provider: {provider}")

            obj, errors, should_continue = _parse_json_response(
                raw, retries, errors, f"{provider}/{model}"
            )
            if should_continue:
                continue
            if obj is None and errors > retries:
                return None
            if obj is not None:
                return obj

        except _RateLimitedError as e:
            pool.mark_rate_limited(e.key, headers=e.headers)
            continue

        except (_anthropic.APIError, ValueError, OSError) as e:
            errors += 1
            if errors > retries:
                _log.warning("LLM call failed (%s/%s): %s", provider, model, e)
                return None

    return None


def call_gemini_embeddings(
    texts: list[str],
    model: str = "gemini-embedding-001",
) -> list | None:
    """Embed texts via Gemini embedContent API with parallel key-rotating calls.

    Uses individual embedContent requests (not batchEmbedContents) so it works
    for all models including experimental ones. Parallelised with ThreadPoolExecutor
    using pool_size * 2 workers. Returns list of float lists or None on total failure.
    """
    import threading  # noqa: PLC0415
    from concurrent.futures import ThreadPoolExecutor  # noqa: PLC0415

    results: list = [None] * len(texts)
    lock = threading.Lock()
    n_workers = max(_gemini_pool.size * 2, 4)

    def _one(args: tuple) -> None:
        idx, text = args
        for _attempt in range(_gemini_pool.size + 3):
            key = _gemini_pool.get_key()
            if not key:
                return
            try:
                url = (
                    f"https://generativelanguage.googleapis.com/v1beta"
                    f"/models/{model}:embedContent"
                )
                resp = _requests.post(
                    url,
                    params={"key": key},
                    json={"content": {"parts": [{"text": text}]}},
                    timeout=60,
                )
                if resp.status_code == _HTTP_TOO_MANY_REQUESTS:
                    _gemini_pool.mark_rate_limited(key, headers={})
                    continue
                resp.raise_for_status()
                values = resp.json().get("embedding", {}).get("values")
                if values:
                    _gemini_pool.clear_cooldown(key)
                    with lock:
                        results[idx] = values
                return
            except _RateLimitedError as e:
                _gemini_pool.mark_rate_limited(e.key, headers=e.headers)
                continue
            except (_requests.RequestException, OSError) as e:
                _log.warning("Gemini embed idx %d failed: %s", idx, e)
                return

    with ThreadPoolExecutor(max_workers=n_workers) as ex:
        list(ex.map(_one, enumerate(texts)))

    return results


def call_openai_embeddings(
    texts: list[str],
    model: str = "text-embedding-3-large",
    batch_size: int = 100,
) -> list | None:
    """Embed a list of texts via OpenAI embeddings API with key rotation.

    Uses _openai_pool for 429 handling. Processes in batches of batch_size.
    Returns list of float lists (3072-dim for text-embedding-3-large), or None on failure.
    """
    results: list = [None] * len(texts)
    max_errors = 3
    errors = 0

    for start in range(0, len(texts), batch_size):
        chunk = texts[start : start + batch_size]
        success = False
        for _attempt in range(_openai_pool.size + 3):
            key = _openai_pool.get_key()
            if not key:
                raise ValueError("No OpenAI API keys configured")
            try:
                client = _get_openai_client(key)
                resp = client.embeddings.create(model=model, input=chunk)
                _openai_pool.clear_cooldown(key)
                for item in resp.data:
                    results[start + item.index] = item.embedding
                success = True
                break
            except _openai.RateLimitError as e:
                headers = {}
                if hasattr(e, "response") and e.response is not None:
                    for name in ("retry-after", "x-ratelimit-reset-requests"):
                        if name in e.response.headers:
                            headers[name] = e.response.headers[name]
                _openai_pool.mark_rate_limited(key, headers=headers)
                continue
            except (_openai.APIError, OSError) as e:
                errors += 1
                _log.warning("Embeddings batch %d failed: %s", start, e)
                if errors > max_errors:
                    return None
                break
        if not success and errors > max_errors:
            return None

    return results


def call_llm_vision_json(  # noqa: C901, PLR0912
    prompt_cfg: dict,
    system_message: str,
    user_message: str,
    image_paths: list[Path],
    provider: str | None = None,
    model: str | None = None,
) -> dict | None:
    """Call an LLM with vision and return parsed JSON dict, or None on failure.

    Each JPEG in *image_paths* is base64-encoded and sent alongside the text.
    Provider and model come from LLM_PROVIDER / LLM_MODEL env vars by default.
    Pass *provider* and *model* to override for this call only.

    Args:
        prompt_cfg: dict with keys temperature (float), max_tokens (int),
            and optional retry_on_invalid_json (int, default 2).
        system_message: System prompt text.
        user_message: User prompt text (appended after image blocks).
        image_paths: List of Path objects pointing to JPEG files.
        provider: LLM provider override ("openai", "gemini", "anthropic"). Defaults to LLM_PROVIDER.
        model: Model name override. Defaults to LLM_MODEL.

    Returns:
        Parsed JSON dict, or None on repeated failure or rate-limit exhaustion.
    """
    provider = provider or LLM_PROVIDER
    model = model or LLM_MODEL
    retries = prompt_cfg.get("retry_on_invalid_json", 2)
    temperature = prompt_cfg["temperature"]
    max_tokens = prompt_cfg["max_tokens"]

    pool = {
        "openai": _openai_pool,
        "gemini": _gemini_pool,
        "anthropic": _anthropic_pool,
    }.get(provider)
    pool_sz = pool.size if pool else 1
    max_attempts = (retries + 1) + pool_sz
    errors = 0

    for _attempt in range(max_attempts):
        try:
            if provider == "openai":
                key = _openai_pool.get_key()
                if not key:
                    raise ValueError("No OpenAI API keys configured")
                raw = _call_openai_vision(
                    key,
                    model,
                    system_message,
                    user_message,
                    temperature,
                    max_tokens,
                    image_paths,
                )
                _openai_pool.clear_cooldown(key)
            elif provider == "gemini":
                key = _gemini_pool.get_key()
                if not key:
                    raise ValueError("No Gemini API keys configured")
                output_tokens = max(max_tokens, GEMINI_MIN_OUTPUT_TOKENS)
                raw = _call_gemini_vision(
                    key,
                    model,
                    system_message,
                    user_message,
                    temperature,
                    output_tokens,
                    image_paths,
                )
                _gemini_pool.clear_cooldown(key)
            elif provider == "anthropic":
                key = _anthropic_pool.get_key()
                if not key:
                    raise ValueError("No Anthropic API keys configured")
                raw = _call_anthropic_vision(
                    key,
                    model,
                    system_message,
                    user_message,
                    temperature,
                    max_tokens,
                    image_paths,
                )
                _anthropic_pool.clear_cooldown(key)
            else:
                raise ValueError(f"Unknown LLM provider: {provider}")

            obj, errors, should_continue = _parse_json_response(
                raw, retries, errors, f"{provider}/{model}"
            )
            if should_continue:
                continue
            if obj is None and errors > retries:
                return None
            if obj is not None:
                return obj

        except _RateLimitedError as e:
            pool.mark_rate_limited(e.key, headers=e.headers)
            continue

        except (_anthropic.APIError, _openai.APIError, ValueError, OSError) as e:
            errors += 1
            if errors > retries:
                _log.warning("LLM vision call failed (%s/%s): %s", provider, model, e)
                return None

    return None
