Saltar a contenido

ciel.providers — Proveedores de modelos

Contratos y proveedores de modelos agnósticos al proveedor. Define los DTOs de configuración (ProviderConfig, ModelInfo), el contrato ChatProvider y un registro/fábrica (ProviderRegistry, ProviderFactory), junto con implementaciones OpenAICompatibleProvider, AnthropicProvider, AzureOpenAIProvider y LiteLLMProvider.

Note

Los DTOs de solicitud/respuesta de chat (ChatRequest, ChatResponse, ChatMessage, ChatChoice) se definen en ciel.runtime y se reutilizan aquí.

Contenido multimodal (ChatMessage.content)

Desde v0.10.0, ChatMessage.content es str | list[dict[str, Any]]:

  • str — texto plano (comportamiento histórico, 100% compatible).
  • list[dict] — lista de partes de contenido (multimodal). Cada parte es un dict con "type" y campos específicos:
type Campos
"text" {"type": "text", "text": "..."}
"image_url" {"type": "image_url", "image_url": {"url": "data:image/png;base64,..."}}
"input_audio" {"type": "input_audio", "input_audio": {"data": "...", "format": "wav"}}

No necesitas mapear estas partes al formato de cada proveedor: los serializers automáticos de cada provider (OpenAICompatibleProvider, AnthropicProvider, GeminiProvider, AzureOpenAIProvider, LiteLLMProvider) convierten el content multimodal a la representación nativa (bloques image de Anthropic, inline_data de Gemini, etc.) de forma transparente.

Método de ayuda:

  • ChatMessage.text() -> str — devuelve el texto plano concatenando las partes de tipo "text" e ignorando imágenes/audio. Útil para CLI, gateway y resúmenes sin romper la API pública.
from ciel.runtime import ChatMessage

# Texto plano (sigue funcionando igual)
msg = ChatMessage(role="user", content="hola")
assert msg.text() == "hola"

# Multimodal: pregunta sobre una imagen (data-URL)
img = ChatMessage(
    role="user",
    content=[
        {"type": "text", "text": "¿Qué hay en esta imagen?"},
        {"type": "image_url", "image_url": {"url": "data:image/png;base64,iVBORw0KGgo..."}},
    ],
)
# .text() solo entrega el texto legible
assert img.text() == "¿Qué hay en esta imagen?"

Proveedores disponibles

Clase Módulo Notas
OpenAICompatibleProvider ciel.providers OpenAI y cualquier endpoint compatible (/v1).
AnthropicProvider ciel.providers Claude. Mapea image_url a bloques image base64.
GeminiProvider ciel.providers.gemini Gemini. Mapea image_url a inline_data.
AzureOpenAIProvider ciel.providers.azure Azure OpenAI (deployment + api-version).
LiteLLMProvider ciel.providers.litellm Meta-provider de 100+ modelos. Requiere el extra litellm.
MockProvider ciel.providers.mock Proveedor determinista offline (modos echo/map/fixed). Sin red.

LiteLLMProvider (extra litellm)

Expone 100+ modelos tras un único contrato ChatProvider delegando en la librería litellm. Es offline-safe: litellm solo se importa al construir/usar el provider, así el import por defecto del framework no arrastra la dependencia pesada.

pip install "mana-ciel[litellm]"
from ciel.providers.litellm import LiteLLMProvider

provider = LiteLLMProvider(
    model="gpt-4o-mini",
    api_key="sk-...",
    # models=[...]  -> opcional: Router LiteLLM para fallback/balanceo
)

Si se omite el extra, la construcción lanza ProviderError claro: The 'litellm' extra is required for LiteLLMProvider. Install it with: pip install "mana-ciel[litellm]".

AzureOpenAIProvider

Azure OpenAI es compatible con OpenAI pero requiere el parámetro de consulta api-version y direcciona un deployment (no el id de modelo crudo). Esta subclase de OpenAICompatibleProvider inyecta la versión en cada request y trata el campo model/deployment como el nombre del deployment.

from ciel.providers.azure import AzureOpenAIProvider

provider = AzureOpenAIProvider(
    base_url="https://<resource>.openai.azure.com",
    api_key="...",            # AZURE_OPENAI_API_KEY
    deployment="gpt-4o",      # nombre del deployment en Azure
    api_version="2024-06-01",
)

Auto-provider por prefijo de modelo

ciel.providers.auto.auto_provider(model) infiere el provider a partir del prefijo del id de modelo (lee la API key del entorno):

Prefijo Provider Base URL por defecto
gpt-, o1, o3 OpenAI-compatible https://api.openai.com/v1
claude- Anthropic https://api.anthropic.com/v1
gemini-, models/ Gemini
azure/ AzureOpenAIProvider AZURE_OPENAI_ENDPOINT
ollama/ OpenAI-compatible (local) http://localhost:11434/v1
vllm/ OpenAI-compatible (self-host) http://localhost:8000/v1 (o VLLM_BASE_URL)
mock/ MockProvider Offline (sin red). Modos echo/map/fixed.

ciel.Agent(model="gpt-4o-mini") usa esto internamente; un provider= explícito siempre gana sobre la inferencia.

ciel.providers

_domain_error = ProviderError module-attribute

AnthropicProvider

Bases: ChatProvider

Source code in src/ciel/providers/__init__.py
class AnthropicProvider(ChatProvider):
    provider_name = "anthropic"

    def __init__(
        self,
        api_key: Optional[str] = None,
        default_model: str = "claude-3-5-haiku-20241022",
        tenant: Optional[str] = None,
        base_url: str = "https://api.anthropic.com/v1",
        timeout: float = 30.0,
    ) -> None:
        self.api_key = api_key
        self.default_model = default_model
        self.tenant = tenant
        self.base_url = base_url.rstrip("/")
        self.timeout = timeout

    async def complete(self, request: "ChatRequest") -> "ChatResponse":
        if self.api_key is None:
            raise _domain_error("Anthropic provider requires api_key")
        model = request.model or self.default_model
        headers = {
            "Content-Type": "application/json",
            "Authorization": f"Bearer {self.api_key}",
            "anthropic-version": "2023-06-01",
        }
        messages = [
            {"role": message.role, "content": _anthropic_content(message.content)}
            for message in request.messages
        ]
        payload = {
            "model": model,
            "messages": messages,
            "max_tokens": request.max_tokens or 64,
        }
        if request.temperature is not None:
            payload["temperature"] = request.temperature
        async with httpx.AsyncClient(base_url=self.base_url, timeout=self.timeout) as client:
            response = await client.post("/messages", headers=headers, json=payload)
            response.raise_for_status()
            body = response.json()
        content = self._extract_text(body)
        choice = _ChatResponse(
            message=_ChatMessage(role="assistant", content=content, metadata={}),
            finish_reason="stop",
            metadata=body,
        )
        from ciel.runtime import ChatResponse as ChatResponseDTO
        return ChatResponseDTO(choice=choice, metadata={"provider": self.provider_name})

    async def stream(self, request: "ChatRequest") -> Sequence["ChatResponse"]:
        return [await self.complete(request)]

    async def models(self) -> Sequence["ModelInfo"]:
        return [ModelInfo(id=self.default_model, provider=self.provider_name, metadata={"tenant": self.tenant})]

    @staticmethod
    def _extract_text(body: Dict[str, Any]) -> str:
        items = body.get("content") or []
        for item in items:
            if isinstance(item, dict) and item.get("type") == "text":
                text = item.get("text", "")
                if isinstance(text, str):
                    return text
        return ""

ChatProvider

Bases: ABC

Source code in src/ciel/providers/__init__.py
class ChatProvider(ABC):
    provider_name: str = "unknown"

    @abstractmethod
    async def complete(self, request: "ChatRequest") -> "ChatResponse":
        raise NotImplementedError

    @abstractmethod
    async def stream(self, request: "ChatRequest") -> Sequence["ChatResponse"]:
        raise NotImplementedError

    @abstractmethod
    async def models(self) -> Sequence[ModelInfo]:
        raise NotImplementedError

CielError

Bases: Exception

Base Ciel error.

Source code in src/ciel/common/__init__.py
class CielError(Exception):
    """Base Ciel error."""

LiteLLMProvider

Bases: ChatProvider

ChatProvider backed by LiteLLM (100+ models, optional Router fallback).

Source code in src/ciel/providers/litellm.py
class LiteLLMProvider(ChatProvider):
    """ChatProvider backed by LiteLLM (100+ models, optional Router fallback)."""

    provider_name = "litellm"

    def __init__(
        self,
        *,
        model: str,
        api_key: Optional[str] = None,
        api_base: Optional[str] = None,
        models: Optional[Sequence[str]] = None,
        tenant: Optional[str] = None,
        timeout: float = 60.0,
        **litellm_kwargs: Any,
    ) -> None:
        # Eagerly import litellm so construction fails clearly when the extra
        # is missing (offline, no network needed — just module availability).
        try:
            self._litellm = _import_litellm()
        except ProviderError:
            raise
        except ImportError as exc:  # pragma: no cover - defensive
            raise ProviderError(
                "The 'litellm' extra is required for LiteLLMProvider. "
                "Install it with: pip install \"mana-ciel[litellm]\""
            ) from exc
        self.model = model
        self.api_key = api_key
        self.api_base = api_base
        self.tenant = tenant
        self.timeout = timeout
        self._litellm_kwargs = dict(litellm_kwargs)
        self._router = None
        if models:
            try:
                Router = self._litellm.Router
                self._router = Router(
                    model_list=[
                        {
                            "model_name": model,
                            "litellm_params": {
                                "model": m,
                                "api_key": api_key,
                                "api_base": api_base,
                                **litellm_kwargs,
                            },
                        }
                        for m in models
                    ]
                )
            except (ImportError, AttributeError):  # pragma: no cover
                self._router = None

    def _common_kwargs(self) -> Dict[str, Any]:
        kwargs: Dict[str, Any] = {"timeout": self.timeout, **self._litellm_kwargs}
        if self.api_key:
            kwargs["api_key"] = self.api_key
        if self.api_base:
            kwargs["api_base"] = self.api_base
        return kwargs

    async def complete(self, request: ChatRequest) -> ChatResponse:
        messages = _build_litellm_messages(request)
        model = request.model or self.model
        kwargs = self._common_kwargs()
        if request.temperature is not None:
            kwargs["temperature"] = request.temperature
        if request.max_tokens is not None:
            kwargs["max_tokens"] = request.max_tokens

        if self._router is not None:
            response = await self._router.acompletion(model=model, messages=messages, **kwargs)
        else:
            response = await self._litellm.acompletion(model=model, messages=messages, **kwargs)

        choice = response.choices[0]
        message = ChatMessage(
            role=choice.message.role or "assistant",
            content=choice.message.content or "",
            name=getattr(choice.message, "name", None),
            metadata={"tenant": self.tenant, "provider": self.provider_name},
        )
        return ChatResponse(
            choice=ChatChoice(
                message=message,
                finish_reason=getattr(choice, "finish_reason", "stop") or "stop",
                usage=getattr(response, "usage", None),
            ),
            metadata={"provider": self.provider_name},
        )

    async def stream(self, request: ChatRequest) -> Sequence[ChatResponse]:
        messages = _build_litellm_messages(request)
        model = request.model or self.model
        kwargs = self._common_kwargs()
        if request.temperature is not None:
            kwargs["temperature"] = request.temperature
        if request.max_tokens is not None:
            kwargs["max_tokens"] = request.max_tokens
        kwargs["stream"] = True

        chunks: List[ChatResponse] = []
        accumulated = ""
        finish_reason = "stop"

        if self._router is not None:
            stream = self._router.acompletion(model=model, messages=messages, **kwargs)
        else:
            stream = self._litellm.acompletion(model=model, messages=messages, **kwargs)
        # LiteLLM returns the async generator directly when stream=True.
        if hasattr(stream, "__await__"):
            stream = await stream

        async for chunk in stream:
            delta = getattr(chunk.choices[0], "delta", None)
            piece = getattr(delta, "content", None) if delta is not None else None
            if isinstance(piece, str) and piece:
                accumulated += piece
            fr = getattr(chunk.choices[0], "finish_reason", None)
            if fr is not None:
                finish_reason = fr or "stop"
            chunks.append(
                ChatResponse(
                    choice=ChatChoice(
                        message=ChatMessage(
                            role="assistant",
                            content=accumulated,
                            metadata={"tenant": self.tenant, "provider": self.provider_name},
                        ),
                        finish_reason=finish_reason,
                    ),
                    metadata={"tenant": self.tenant, "streaming": True, "provider": self.provider_name},
                )
            )
        return tuple(chunks)

    async def models(self) -> Sequence[ModelInfo]:
        if self._router is not None:
            try:
                deployment_names = [d for d in self._router.get_model_names()]  # type: ignore[attr-defined]
            except Exception:  # pragma: no cover - router API variance
                deployment_names = [self.model]
            return [
                ModelInfo(id=name, provider=self.provider_name, metadata={"tenant": self.tenant})
                for name in deployment_names
            ]
        return [ModelInfo(id=self.model, provider=self.provider_name, metadata={"tenant": self.tenant})]

MockProvider

Bases: ChatProvider

Proveedor determinista configurable para eval/tests sin red.

Parameters:

Name Type Description Default
mode str

"echo" | "map" | "fixed".

'fixed'
response str

respuesta constante (modo fixed).

''
mapping Optional[Dict[str, str]]

dict prompt -> response (modo map). La coincidencia es exacta; si no hay exacta, se usa la primera clave que sea substring del prompt (case-insensitive).

None
model str

id de modelo devuelto por models() (default "mock").

'mock'
tenant Optional[str]

tenant_id propagado en metadata.

None
Source code in src/ciel/providers/mock.py
class MockProvider(ChatProvider):
    """Proveedor determinista configurable para eval/tests sin red.

    Args:
        mode: ``"echo"`` | ``"map"`` | ``"fixed"``.
        response: respuesta constante (modo ``fixed``).
        mapping: dict ``prompt -> response`` (modo ``map``). La coincidencia es
            exacta; si no hay exacta, se usa la primera clave que sea substring
            del prompt (case-insensitive).
        model: id de modelo devuelto por ``models()`` (default ``"mock"``).
        tenant: tenant_id propagado en metadata.
    """

    provider_name = "mock"

    def __init__(
        self,
        *,
        mode: str = "fixed",
        response: str = "",
        mapping: Optional[Dict[str, str]] = None,
        model: str = "mock",
        tenant: Optional[str] = None,
    ) -> None:
        if mode not in ("echo", "map", "fixed"):
            raise ValueError(f"MockProvider mode inválido: {mode!r}; usar echo|map|fixed")
        self.mode = mode
        self.response = response
        self.mapping = dict(mapping or {})
        self.model = model
        self.tenant = tenant

    def _respond(self, request: Any) -> str:
        if self.mode == "fixed":
            return self.response
        if self.mode == "echo":
            last = request.messages[-1].content if request.messages else ""
            if isinstance(last, list):
                last = "".join(p.get("text", "") for p in last if isinstance(p, dict))
            last = last or ""
            words = last.split()
            return words[-1] if words else ""
        # map
        prompt = request.messages[-1].content if request.messages else ""
        if isinstance(prompt, list):
            prompt = "".join(p.get("text", "") for p in prompt if isinstance(p, dict))
        prompt = prompt or ""
        if prompt in self.mapping:
            return self.mapping[prompt]
        lowered = prompt.lower()
        for key, val in self.mapping.items():
            if key and key.lower() in lowered:
                return val
        return self.mapping.get("", self.response)

    async def complete(self, request: Any) -> ChatResponse:
        text = self._respond(request)
        message = ChatMessage(role="assistant", content=text, metadata={"tenant": self.tenant})
        return ChatResponse(
            choice=ChatChoice(message=message, finish_reason="stop"),
            metadata={"provider": self.provider_name, "mode": self.mode, "tenant": self.tenant},
        )

    async def stream(self, request: Any) -> Sequence[ChatResponse]:
        return (await self.complete(request),)

    async def models(self) -> Sequence[ModelInfo]:
        return [
            ModelInfo(
                id=self.model,
                provider=self.provider_name,
                metadata={"tenant": self.tenant, "mode": self.mode},
            )
        ]

ModelInfo dataclass

Source code in src/ciel/providers/__init__.py
@dataclass(frozen=True)
class ModelInfo:
    id: str
    provider: str
    capabilities: Sequence[str] = ()
    context_window: Optional[int] = None
    metadata: Dict[str, Any] = field(default_factory=dict)

OpenAICompatibleProvider

Bases: ChatProvider

Source code in src/ciel/providers/__init__.py
class OpenAICompatibleProvider(ChatProvider):
    provider_name = "openai_compat"

    def __init__(
        self,
        *,
        base_url: str,
        api_key: Optional[str] = None,
        default_model: Optional[str] = None,
        timeout: float = 30.0,
        tenant: Optional[str] = None,
    ) -> None:
        self.base_url = base_url.rstrip("/")
        self.api_key = api_key
        self.default_model = default_model
        self.timeout = timeout
        self.tenant = tenant

    def _client_ctx(self):
        """Context manager yielding an ``httpx.AsyncClient`` for this provider."""
        return httpx.AsyncClient(base_url=self.base_url, timeout=self.timeout)

    async def complete(self, request: "ChatRequest") -> "ChatResponse":
        model = request.model or self.default_model or "unknown"
        headers: Dict[str, str] = {"Content-Type": "application/json"}
        if self.api_key:
            headers["Authorization"] = f"Bearer {self.api_key}"

        payload = {
            "model": model,
            "messages": [_serialize_chat_message(asdict(message)) for message in request.messages],
            "temperature": request.temperature,
            "max_tokens": request.max_tokens,
        }
        payload = {key: value for key, value in payload.items() if value is not None}

        async with httpx.AsyncClient(base_url=self.base_url, timeout=self.timeout) as client:
            response = await client.post("/chat/completions", headers=headers, json=payload)
            response.raise_for_status()
            body = response.json()
        from ciel.runtime import ChatChoice, ChatMessage, ChatResponse
        choice = body["choices"][0]
        message = ChatMessage(
            role=choice["message"]["role"],
            content=choice["message"].get("content", ""),
            name=choice["message"].get("name"),
            tool_call_id=choice["message"].get("tool_call_id"),
            metadata=choice["message"].get("metadata", {}),
        )
        return ChatResponse(
            choice=ChatChoice(message=message, finish_reason=choice.get("finish_reason", "stop")),
            metadata=body.get("usage", {}),
        )

    async def stream(self, request: "ChatRequest") -> Sequence["ChatResponse"]:
        from ciel.runtime import ChatChoice, ChatMessage, ChatResponse

        model = request.model or self.default_model or "unknown"
        headers: Dict[str, str] = {"Content-Type": "application/json", "Accept": "text/event-stream"}
        if self.api_key:
            headers["Authorization"] = f"Bearer {self.api_key}"

        payload = {
            "model": model,
            "messages": [_serialize_chat_message(asdict(message)) for message in request.messages],
            "temperature": request.temperature,
            "max_tokens": request.max_tokens,
            "stream": True,
        }
        payload = {key: value for key, value in payload.items() if value is not None}

        chunks: list["ChatResponse"] = []
        accumulated = ""
        finish_reason = "stop"

        async with httpx.AsyncClient(base_url=self.base_url, timeout=self.timeout) as client:
            async with client.stream("POST", "/chat/completions", headers=headers, json=payload) as response:
                response.raise_for_status()
                async for raw_line in response.aiter_lines():
                    line = raw_line.strip()
                    if not line or not line.startswith("data:"):
                        continue
                    data = line[len("data:"):].strip()
                    if data == "[DONE]":
                        break
                    try:
                        event = json.loads(data)
                    except json.JSONDecodeError:
                        continue
                    if not isinstance(event, dict):
                        continue
                    choices = event.get("choices") or []
                    if not choices:
                        continue
                    delta = choices[0].get("delta") or {}
                    piece = delta.get("content")
                    if isinstance(piece, str) and piece:
                        accumulated += piece
                    if choices[0].get("finish_reason") is not None:
                        finish_reason = choices[0].get("finish_reason") or "stop"
                    message = ChatMessage(
                        role="assistant",
                        content=accumulated,
                        metadata={"tenant": self.tenant},
                    )
                    chunks.append(
                        ChatResponse(
                            choice=ChatChoice(message=message, finish_reason=finish_reason),
                            metadata={"tenant": self.tenant, "streaming": True},
                        )
                    )
        return tuple(chunks)

    async def models(self) -> Sequence[ModelInfo]:
        headers: Dict[str, str] = {}
        if self.api_key:
            headers["Authorization"] = f"Bearer {self.api_key}"
        try:
            async with httpx.AsyncClient(base_url=self.base_url, timeout=self.timeout) as client:
                response = await client.get("/models", headers=headers)
                response.raise_for_status()
                body = response.json()
            model_ids = [item.get("id") for item in body.get("data", []) if item.get("id")]
            if model_ids:
                return [ModelInfo(id=item, provider=self.provider_name, metadata={"tenant": self.tenant}) for item in model_ids]
        except Exception as exc:  # pragma: no cover - network failure path
            raise _domain_error(f"Failed to list models: {exc}") from exc
        return [ModelInfo(id=self.default_model or "unknown", provider=self.provider_name, metadata={"tenant": self.tenant})]

ProviderConfig dataclass

Source code in src/ciel/providers/__init__.py
@dataclass(frozen=True)
class ProviderConfig:
    name: str
    base_url: str
    api_key: Optional[str] = None
    default_model: Optional[str] = None
    timeout: float = 30.0
    tenant: Optional[str] = None

ProviderError

Bases: CielError

Model/provider error.

Source code in src/ciel/common/__init__.py
class ProviderError(CielError):
    """Model/provider error."""

ProviderFactory

Source code in src/ciel/providers/__init__.py
class ProviderFactory:
    @staticmethod
    def from_config(config: ProviderConfig) -> ChatProvider:
        normalized = config.base_url.rstrip("/")
        if config.name == "litellm":
            from ciel.providers.litellm import LiteLLMProvider

            return LiteLLMProvider(
                model=config.default_model or "gpt-4o-mini",
                api_key=config.api_key,
                api_base=config.base_url if config.base_url else None,
                tenant=config.tenant,
                timeout=config.timeout,
            )
        if normalized.endswith("/v1") or "openai" in normalized:
            return OpenAICompatibleProvider(
                base_url=normalized,
                api_key=config.api_key,
                default_model=config.default_model,
                timeout=config.timeout,
                tenant=config.tenant,
            )
        raise ProviderError(f"No provider implementation for base_url: {config.base_url}")

ProviderRegistry

Source code in src/ciel/providers/__init__.py
class ProviderRegistry:
    def __init__(self) -> None:
        self._providers: Dict[str, ChatProvider] = {}
        self._configs: Dict[str, Dict[str, Any]] = {}

    def register(self, name: str, provider: ChatProvider, *, config: Optional[Dict[str, Any]] = None) -> None:
        self._providers[name] = provider
        self._configs[name] = config or {}

    def get(self, name: str) -> ChatProvider:
        if name not in self._providers:
            raise _domain_error(f"Provider not registered: {name}")
        return self._providers[name]

    def available(self) -> Sequence[str]:
        return list(self._providers.keys())

_ChatMessage dataclass

Source code in src/ciel/providers/__init__.py
@dataclass
class _ChatMessage:
    role: str
    content: str
    name: Optional[str] = None
    metadata: Dict[str, Any] = field(default_factory=dict)

_ChatResponse dataclass

Source code in src/ciel/providers/__init__.py
@dataclass
class _ChatResponse:
    message: _ChatMessage
    finish_reason: str = "stop"
    metadata: Dict[str, Any] = field(default_factory=dict)

_anthropic_content(content: Any) -> List[Dict[str, Any]]

Normalize str | list[dict] content to Anthropic content blocks.

  • str -> [{"type": "text", "text": <content>}]
  • list -> maps each part to Anthropic blocks:
    • {"type": "text", "text": ...} -> text block
    • {"type": "image_url", "image_url": {"url": "data:..."}} -> {"type": "image", "source": {"type": "base64", "media_type": ..., "data": <b64>}} (only data-URLs are supported by the API)
Source code in src/ciel/providers/__init__.py
def _anthropic_content(content: Any) -> List[Dict[str, Any]]:
    """Normalize ``str | list[dict]`` content to Anthropic content blocks.

    - ``str`` -> ``[{"type": "text", "text": <content>}]``
    - ``list`` -> maps each part to Anthropic blocks:
        * ``{"type": "text", "text": ...}`` -> text block
        * ``{"type": "image_url", "image_url": {"url": "data:..."}}`` ->
          ``{"type": "image", "source": {"type": "base64", "media_type": ...,
          "data": <b64>}}`` (only data-URLs are supported by the API)
    """
    if isinstance(content, str):
        return [{"type": "text", "text": content}]
    if not isinstance(content, list):
        return [{"type": "text", "text": ""}]
    blocks: List[Dict[str, Any]] = []
    for part in content:
        if not isinstance(part, dict):
            continue
        ptype = part.get("type")
        if ptype == "text":
            text = part.get("text", "")
            if isinstance(text, str):
                blocks.append({"type": "text", "text": text})
        elif ptype == "image_url":
            url = (part.get("image_url") or {}).get("url", "")
            media_type, data = _split_data_url(url)
            if media_type is not None and data is not None:
                blocks.append(
                    {
                        "type": "image",
                        "source": {"type": "base64", "media_type": media_type, "data": data},
                    }
                )
    return blocks or [{"type": "text", "text": ""}]

_build_chat_message(payload: Dict[str, Any]) -> _ChatMessage

Source code in src/ciel/providers/__init__.py
def _build_chat_message(payload: Dict[str, Any]) -> _ChatMessage:
    return _ChatMessage(
        role=payload["role"],
        content=payload.get("content", "") or "",
        name=payload.get("name"),
        metadata=payload.get("metadata", {}),
    )

_normalize_content_parts(content: Any) -> List[Dict[str, Any]]

Normalize str | list[dict] content to OpenAI content-parts.

  • str -> [{"type": "text", "text": <content>}]
  • list -> validated list of part dicts (text/image_url/input_audio/...), tolerating a bare {"text": ...} part by promoting it to a text part.
Source code in src/ciel/providers/__init__.py
def _normalize_content_parts(content: Any) -> List[Dict[str, Any]]:
    """Normalize ``str | list[dict]`` content to OpenAI content-parts.

    - ``str`` -> ``[{"type": "text", "text": <content>}]``
    - ``list`` -> validated list of part dicts (text/image_url/input_audio/...),
      tolerating a bare ``{"text": ...}`` part by promoting it to a text part.
    """
    if isinstance(content, str):
        return [{"type": "text", "text": content}]
    if isinstance(content, list):
        parts: List[Dict[str, Any]] = []
        for part in content:
            if isinstance(part, dict):
                if "type" not in part and "text" in part:
                    part = {"type": "text", "text": part["text"]}
                parts.append(part)
        return parts
    return [{"type": "text", "text": ""}]

_serialize_chat_message(message: Dict[str, Any]) -> Dict[str, Any]

Source code in src/ciel/providers/__init__.py
def _serialize_chat_message(message: Dict[str, Any]) -> Dict[str, Any]:
    content = message.get("content", "")
    payload: Dict[str, Any] = {"role": message["role"], "content": _normalize_content_parts(content)}
    name = message.get("name")
    if name:
        payload["name"] = name
    return payload

_split_data_url(url: str) -> Tuple[Optional[str], Optional[str]]

Split a data:<mime>;base64,<data> URL into (mime, b64).

Source code in src/ciel/providers/__init__.py
def _split_data_url(url: str) -> Tuple[Optional[str], Optional[str]]:
    """Split a ``data:<mime>;base64,<data>`` URL into (mime, b64)."""
    if not isinstance(url, str) or not url.startswith("data:"):
        return None, None
    try:
        meta, data = url[5:].split(",", 1)
    except ValueError:
        return None, None
    if ";base64" not in meta:
        return None, None
    mime = meta.split(";base64", 1)[0] or "application/octet-stream"
    return mime, data

Azure OpenAI provider (Fase 16-B).

Azure OpenAI is OpenAI-compatible but requires the api-version query parameter and addresses a deployment name rather than the raw model id. This subclass injects the api-version on every request and treats the model field as the deployment name. Offline-safe: no network at construction.

AzureOpenAIProvider

Bases: OpenAICompatibleProvider

Source code in src/ciel/providers/azure.py
class AzureOpenAIProvider(OpenAICompatibleProvider):
    provider_name = "azure_openai"

    def __init__(
        self,
        *,
        base_url: str,
        api_key: Optional[str] = None,
        api_version: str = "2024-06-01",
        deployment: Optional[str] = None,
        default_model: Optional[str] = None,
        tenant: Optional[str] = None,
        timeout: float = 30.0,
    ) -> None:
        # Azure base_url looks like https://<resource>.openai.azure.com
        super().__init__(
            base_url=base_url.rstrip("/"),
            api_key=api_key,
            default_model=deployment or default_model,
            timeout=timeout,
            tenant=tenant,
        )
        self.api_version = api_version
        self.deployment = deployment or default_model

    def _chat_path(self, model: str) -> str:
        deployment = self.deployment or model
        return f"/openai/deployments/{deployment}/chat/completions?api-version={self.api_version}"

    def _build_messages(self, request: ChatRequest) -> list[Dict[str, Any]]:
        return [
            {
                "role": m.role,
                "content": _normalize_content_parts(m.content),
                **({"name": m.name} if m.name else {}),
            }
            for m in request.messages
        ]

    async def complete(self, request: ChatRequest) -> ChatResponse:
        model = request.model or self.default_model or (self.deployment or "unknown")
        headers: Dict[str, str] = {"Content-Type": "application/json"}
        if self.api_key:
            headers["api-key"] = self.api_key

        payload: Dict[str, Any] = {
            "model": model,
            "messages": self._build_messages(request),
            "temperature": request.temperature,
            "max_tokens": request.max_tokens,
        }
        payload = {k: v for k, v in payload.items() if v is not None}

        async with self._client_ctx() as client:
            response = await client.post(self._chat_path(model), headers=headers, json=payload)
            response.raise_for_status()
            body = response.json()
        choice = body["choices"][0]
        message = ChatMessage(
            role=choice["message"].get("role", "assistant"),
            content=choice["message"].get("content", ""),
            name=choice["message"].get("name"),
            tool_call_id=choice["message"].get("tool_call_id"),
            metadata={"tenant": self.tenant, "provider": self.provider_name},
        )
        return ChatResponse(
            choice=ChatChoice(message=message, finish_reason=choice.get("finish_reason", "stop")),
            metadata=body.get("usage", {}),
        )

    async def stream(self, request: ChatRequest) -> Sequence[ChatResponse]:
        model = request.model or self.default_model or (self.deployment or "unknown")
        headers: Dict[str, str] = {"Content-Type": "application/json", "Accept": "text/event-stream"}
        if self.api_key:
            headers["api-key"] = self.api_key

        payload: Dict[str, Any] = {
            "model": model,
            "messages": self._build_messages(request),
            "temperature": request.temperature,
            "max_tokens": request.max_tokens,
            "stream": True,
        }
        payload = {k: v for k, v in payload.items() if v is not None}

        chunks: list[ChatResponse] = []
        accumulated = ""
        finish_reason = "stop"
        async with self._client_ctx() as client:
            async with client.stream("POST", self._chat_path(model), headers=headers, json=payload) as response:
                response.raise_for_status()
                async for raw_line in response.aiter_lines():
                    line = raw_line.strip()
                    if not line or not line.startswith("data:"):
                        continue
                    data = line[len("data:"):].strip()
                    if data == "[DONE]":
                        break
                    try:
                        event = json.loads(data)
                    except json.JSONDecodeError:
                        continue
                    if not isinstance(event, dict):
                        continue
                    choices = event.get("choices") or []
                    if not choices:
                        continue
                    delta = choices[0].get("delta") or {}
                    piece = delta.get("content")
                    if isinstance(piece, str) and piece:
                        accumulated += piece
                    if choices[0].get("finish_reason") is not None:
                        finish_reason = choices[0].get("finish_reason") or "stop"
                    chunks.append(
                        ChatResponse(
                            choice=ChatChoice(
                                message=ChatMessage(
                                    role="assistant",
                                    content=accumulated,
                                    metadata={"tenant": self.tenant, "provider": self.provider_name},
                                ),
                                finish_reason=finish_reason,
                            ),
                            metadata={"tenant": self.tenant, "streaming": True, "provider": self.provider_name},
                        )
                    )
        return tuple(chunks)

    async def models(self) -> Sequence[ModelInfo]:
        name = self.deployment or self.default_model or "unknown"
        return [ModelInfo(id=name, provider=self.provider_name, metadata={"tenant": self.tenant})]

ChatChoice dataclass

Source code in src/ciel/runtime/tools.py
@dataclass(frozen=True)
class ChatChoice:
    message: ChatMessage
    finish_reason: str
    usage: Optional[Dict[str, Any]] = None
    metadata: Dict[str, Any] = field(default_factory=dict)

ChatMessage dataclass

Source code in src/ciel/runtime/tools.py
@dataclass(frozen=True)
class ChatMessage:
    role: str
    content: ChatContent
    name: Optional[str] = None
    tool_call_id: Optional[str] = None
    tool_calls: Optional[list[dict[str, Any]]] = None
    metadata: Dict[str, Any] = field(default_factory=dict)

    def text(self) -> str:
        """Extract plain text from content, tolerant to multimodal parts.

        - ``str`` content is returned verbatim.
        - ``list`` content concatenates the ``text`` of every part whose type
          is ``"text"`` but drops images/audio/video, so consumers (CLI,
          compression, ``AgentResponse.text``) see only readable text.
        """
        content = self.content
        if isinstance(content, str):
            return content
        if not isinstance(content, list):
            return ""
        parts: list[str] = []
        for part in content:
            if isinstance(part, dict) and part.get("type") == "text":
                text = part.get("text", "")
                if isinstance(text, str):
                    parts.append(text)
        return "".join(parts)

text() -> str

Extract plain text from content, tolerant to multimodal parts.

  • str content is returned verbatim.
  • list content concatenates the text of every part whose type is "text" but drops images/audio/video, so consumers (CLI, compression, AgentResponse.text) see only readable text.
Source code in src/ciel/runtime/tools.py
def text(self) -> str:
    """Extract plain text from content, tolerant to multimodal parts.

    - ``str`` content is returned verbatim.
    - ``list`` content concatenates the ``text`` of every part whose type
      is ``"text"`` but drops images/audio/video, so consumers (CLI,
      compression, ``AgentResponse.text``) see only readable text.
    """
    content = self.content
    if isinstance(content, str):
        return content
    if not isinstance(content, list):
        return ""
    parts: list[str] = []
    for part in content:
        if isinstance(part, dict) and part.get("type") == "text":
            text = part.get("text", "")
            if isinstance(text, str):
                parts.append(text)
    return "".join(parts)

ChatProvider

Bases: ABC

Source code in src/ciel/providers/__init__.py
class ChatProvider(ABC):
    provider_name: str = "unknown"

    @abstractmethod
    async def complete(self, request: "ChatRequest") -> "ChatResponse":
        raise NotImplementedError

    @abstractmethod
    async def stream(self, request: "ChatRequest") -> Sequence["ChatResponse"]:
        raise NotImplementedError

    @abstractmethod
    async def models(self) -> Sequence[ModelInfo]:
        raise NotImplementedError

ChatRequest dataclass

Source code in src/ciel/runtime/tools.py
@dataclass(frozen=True)
class ChatRequest:
    messages: Sequence[ChatMessage]
    tools: Sequence[ToolSpec] = ()
    model: Optional[str] = None
    temperature: Optional[float] = None
    max_tokens: Optional[int] = None
    extra: Dict[str, Any] = field(default_factory=dict)

ChatResponse dataclass

Source code in src/ciel/runtime/tools.py
@dataclass(frozen=True)
class ChatResponse:
    choice: ChatChoice
    metadata: Dict[str, Any] = field(default_factory=dict)

ModelInfo dataclass

Source code in src/ciel/providers/__init__.py
@dataclass(frozen=True)
class ModelInfo:
    id: str
    provider: str
    capabilities: Sequence[str] = ()
    context_window: Optional[int] = None
    metadata: Dict[str, Any] = field(default_factory=dict)

OpenAICompatibleProvider

Bases: ChatProvider

Source code in src/ciel/providers/__init__.py
class OpenAICompatibleProvider(ChatProvider):
    provider_name = "openai_compat"

    def __init__(
        self,
        *,
        base_url: str,
        api_key: Optional[str] = None,
        default_model: Optional[str] = None,
        timeout: float = 30.0,
        tenant: Optional[str] = None,
    ) -> None:
        self.base_url = base_url.rstrip("/")
        self.api_key = api_key
        self.default_model = default_model
        self.timeout = timeout
        self.tenant = tenant

    def _client_ctx(self):
        """Context manager yielding an ``httpx.AsyncClient`` for this provider."""
        return httpx.AsyncClient(base_url=self.base_url, timeout=self.timeout)

    async def complete(self, request: "ChatRequest") -> "ChatResponse":
        model = request.model or self.default_model or "unknown"
        headers: Dict[str, str] = {"Content-Type": "application/json"}
        if self.api_key:
            headers["Authorization"] = f"Bearer {self.api_key}"

        payload = {
            "model": model,
            "messages": [_serialize_chat_message(asdict(message)) for message in request.messages],
            "temperature": request.temperature,
            "max_tokens": request.max_tokens,
        }
        payload = {key: value for key, value in payload.items() if value is not None}

        async with httpx.AsyncClient(base_url=self.base_url, timeout=self.timeout) as client:
            response = await client.post("/chat/completions", headers=headers, json=payload)
            response.raise_for_status()
            body = response.json()
        from ciel.runtime import ChatChoice, ChatMessage, ChatResponse
        choice = body["choices"][0]
        message = ChatMessage(
            role=choice["message"]["role"],
            content=choice["message"].get("content", ""),
            name=choice["message"].get("name"),
            tool_call_id=choice["message"].get("tool_call_id"),
            metadata=choice["message"].get("metadata", {}),
        )
        return ChatResponse(
            choice=ChatChoice(message=message, finish_reason=choice.get("finish_reason", "stop")),
            metadata=body.get("usage", {}),
        )

    async def stream(self, request: "ChatRequest") -> Sequence["ChatResponse"]:
        from ciel.runtime import ChatChoice, ChatMessage, ChatResponse

        model = request.model or self.default_model or "unknown"
        headers: Dict[str, str] = {"Content-Type": "application/json", "Accept": "text/event-stream"}
        if self.api_key:
            headers["Authorization"] = f"Bearer {self.api_key}"

        payload = {
            "model": model,
            "messages": [_serialize_chat_message(asdict(message)) for message in request.messages],
            "temperature": request.temperature,
            "max_tokens": request.max_tokens,
            "stream": True,
        }
        payload = {key: value for key, value in payload.items() if value is not None}

        chunks: list["ChatResponse"] = []
        accumulated = ""
        finish_reason = "stop"

        async with httpx.AsyncClient(base_url=self.base_url, timeout=self.timeout) as client:
            async with client.stream("POST", "/chat/completions", headers=headers, json=payload) as response:
                response.raise_for_status()
                async for raw_line in response.aiter_lines():
                    line = raw_line.strip()
                    if not line or not line.startswith("data:"):
                        continue
                    data = line[len("data:"):].strip()
                    if data == "[DONE]":
                        break
                    try:
                        event = json.loads(data)
                    except json.JSONDecodeError:
                        continue
                    if not isinstance(event, dict):
                        continue
                    choices = event.get("choices") or []
                    if not choices:
                        continue
                    delta = choices[0].get("delta") or {}
                    piece = delta.get("content")
                    if isinstance(piece, str) and piece:
                        accumulated += piece
                    if choices[0].get("finish_reason") is not None:
                        finish_reason = choices[0].get("finish_reason") or "stop"
                    message = ChatMessage(
                        role="assistant",
                        content=accumulated,
                        metadata={"tenant": self.tenant},
                    )
                    chunks.append(
                        ChatResponse(
                            choice=ChatChoice(message=message, finish_reason=finish_reason),
                            metadata={"tenant": self.tenant, "streaming": True},
                        )
                    )
        return tuple(chunks)

    async def models(self) -> Sequence[ModelInfo]:
        headers: Dict[str, str] = {}
        if self.api_key:
            headers["Authorization"] = f"Bearer {self.api_key}"
        try:
            async with httpx.AsyncClient(base_url=self.base_url, timeout=self.timeout) as client:
                response = await client.get("/models", headers=headers)
                response.raise_for_status()
                body = response.json()
            model_ids = [item.get("id") for item in body.get("data", []) if item.get("id")]
            if model_ids:
                return [ModelInfo(id=item, provider=self.provider_name, metadata={"tenant": self.tenant}) for item in model_ids]
        except Exception as exc:  # pragma: no cover - network failure path
            raise _domain_error(f"Failed to list models: {exc}") from exc
        return [ModelInfo(id=self.default_model or "unknown", provider=self.provider_name, metadata={"tenant": self.tenant})]

ProviderError

Bases: CielError

Model/provider error.

Source code in src/ciel/common/__init__.py
class ProviderError(CielError):
    """Model/provider error."""

_normalize_content_parts(content: Any) -> List[Dict[str, Any]]

Normalize str | list[dict] content to OpenAI content-parts.

  • str -> [{"type": "text", "text": <content>}]
  • list -> validated list of part dicts (text/image_url/input_audio/...), tolerating a bare {"text": ...} part by promoting it to a text part.
Source code in src/ciel/providers/__init__.py
def _normalize_content_parts(content: Any) -> List[Dict[str, Any]]:
    """Normalize ``str | list[dict]`` content to OpenAI content-parts.

    - ``str`` -> ``[{"type": "text", "text": <content>}]``
    - ``list`` -> validated list of part dicts (text/image_url/input_audio/...),
      tolerating a bare ``{"text": ...}`` part by promoting it to a text part.
    """
    if isinstance(content, str):
        return [{"type": "text", "text": content}]
    if isinstance(content, list):
        parts: List[Dict[str, Any]] = []
        for part in content:
            if isinstance(part, dict):
                if "type" not in part and "text" in part:
                    part = {"type": "text", "text": part["text"]}
                parts.append(part)
        return parts
    return [{"type": "text", "text": ""}]

LiteLLM meta-provider (Fase 16-A).

Exposes 100+ models through a single :class:ChatProvider contract by delegating to the litellm library. This module is offline-safe: it never imports litellm at module load time, only when the provider is actually constructed/used, so the default framework import graph stays free of the heavy litellm dependency. Install it with the optional extra::

pip install "mana-ciel[litellm]"

Fallback/load-balancing across multiple models is available via a LiteLLM Router (passed as models); when omitted, a single model is used.

ChatChoice dataclass

Source code in src/ciel/runtime/tools.py
@dataclass(frozen=True)
class ChatChoice:
    message: ChatMessage
    finish_reason: str
    usage: Optional[Dict[str, Any]] = None
    metadata: Dict[str, Any] = field(default_factory=dict)

ChatMessage dataclass

Source code in src/ciel/runtime/tools.py
@dataclass(frozen=True)
class ChatMessage:
    role: str
    content: ChatContent
    name: Optional[str] = None
    tool_call_id: Optional[str] = None
    tool_calls: Optional[list[dict[str, Any]]] = None
    metadata: Dict[str, Any] = field(default_factory=dict)

    def text(self) -> str:
        """Extract plain text from content, tolerant to multimodal parts.

        - ``str`` content is returned verbatim.
        - ``list`` content concatenates the ``text`` of every part whose type
          is ``"text"`` but drops images/audio/video, so consumers (CLI,
          compression, ``AgentResponse.text``) see only readable text.
        """
        content = self.content
        if isinstance(content, str):
            return content
        if not isinstance(content, list):
            return ""
        parts: list[str] = []
        for part in content:
            if isinstance(part, dict) and part.get("type") == "text":
                text = part.get("text", "")
                if isinstance(text, str):
                    parts.append(text)
        return "".join(parts)

text() -> str

Extract plain text from content, tolerant to multimodal parts.

  • str content is returned verbatim.
  • list content concatenates the text of every part whose type is "text" but drops images/audio/video, so consumers (CLI, compression, AgentResponse.text) see only readable text.
Source code in src/ciel/runtime/tools.py
def text(self) -> str:
    """Extract plain text from content, tolerant to multimodal parts.

    - ``str`` content is returned verbatim.
    - ``list`` content concatenates the ``text`` of every part whose type
      is ``"text"`` but drops images/audio/video, so consumers (CLI,
      compression, ``AgentResponse.text``) see only readable text.
    """
    content = self.content
    if isinstance(content, str):
        return content
    if not isinstance(content, list):
        return ""
    parts: list[str] = []
    for part in content:
        if isinstance(part, dict) and part.get("type") == "text":
            text = part.get("text", "")
            if isinstance(text, str):
                parts.append(text)
    return "".join(parts)

ChatProvider

Bases: ABC

Source code in src/ciel/providers/__init__.py
class ChatProvider(ABC):
    provider_name: str = "unknown"

    @abstractmethod
    async def complete(self, request: "ChatRequest") -> "ChatResponse":
        raise NotImplementedError

    @abstractmethod
    async def stream(self, request: "ChatRequest") -> Sequence["ChatResponse"]:
        raise NotImplementedError

    @abstractmethod
    async def models(self) -> Sequence[ModelInfo]:
        raise NotImplementedError

ChatRequest dataclass

Source code in src/ciel/runtime/tools.py
@dataclass(frozen=True)
class ChatRequest:
    messages: Sequence[ChatMessage]
    tools: Sequence[ToolSpec] = ()
    model: Optional[str] = None
    temperature: Optional[float] = None
    max_tokens: Optional[int] = None
    extra: Dict[str, Any] = field(default_factory=dict)

ChatResponse dataclass

Source code in src/ciel/runtime/tools.py
@dataclass(frozen=True)
class ChatResponse:
    choice: ChatChoice
    metadata: Dict[str, Any] = field(default_factory=dict)

LiteLLMProvider

Bases: ChatProvider

ChatProvider backed by LiteLLM (100+ models, optional Router fallback).

Source code in src/ciel/providers/litellm.py
class LiteLLMProvider(ChatProvider):
    """ChatProvider backed by LiteLLM (100+ models, optional Router fallback)."""

    provider_name = "litellm"

    def __init__(
        self,
        *,
        model: str,
        api_key: Optional[str] = None,
        api_base: Optional[str] = None,
        models: Optional[Sequence[str]] = None,
        tenant: Optional[str] = None,
        timeout: float = 60.0,
        **litellm_kwargs: Any,
    ) -> None:
        # Eagerly import litellm so construction fails clearly when the extra
        # is missing (offline, no network needed — just module availability).
        try:
            self._litellm = _import_litellm()
        except ProviderError:
            raise
        except ImportError as exc:  # pragma: no cover - defensive
            raise ProviderError(
                "The 'litellm' extra is required for LiteLLMProvider. "
                "Install it with: pip install \"mana-ciel[litellm]\""
            ) from exc
        self.model = model
        self.api_key = api_key
        self.api_base = api_base
        self.tenant = tenant
        self.timeout = timeout
        self._litellm_kwargs = dict(litellm_kwargs)
        self._router = None
        if models:
            try:
                Router = self._litellm.Router
                self._router = Router(
                    model_list=[
                        {
                            "model_name": model,
                            "litellm_params": {
                                "model": m,
                                "api_key": api_key,
                                "api_base": api_base,
                                **litellm_kwargs,
                            },
                        }
                        for m in models
                    ]
                )
            except (ImportError, AttributeError):  # pragma: no cover
                self._router = None

    def _common_kwargs(self) -> Dict[str, Any]:
        kwargs: Dict[str, Any] = {"timeout": self.timeout, **self._litellm_kwargs}
        if self.api_key:
            kwargs["api_key"] = self.api_key
        if self.api_base:
            kwargs["api_base"] = self.api_base
        return kwargs

    async def complete(self, request: ChatRequest) -> ChatResponse:
        messages = _build_litellm_messages(request)
        model = request.model or self.model
        kwargs = self._common_kwargs()
        if request.temperature is not None:
            kwargs["temperature"] = request.temperature
        if request.max_tokens is not None:
            kwargs["max_tokens"] = request.max_tokens

        if self._router is not None:
            response = await self._router.acompletion(model=model, messages=messages, **kwargs)
        else:
            response = await self._litellm.acompletion(model=model, messages=messages, **kwargs)

        choice = response.choices[0]
        message = ChatMessage(
            role=choice.message.role or "assistant",
            content=choice.message.content or "",
            name=getattr(choice.message, "name", None),
            metadata={"tenant": self.tenant, "provider": self.provider_name},
        )
        return ChatResponse(
            choice=ChatChoice(
                message=message,
                finish_reason=getattr(choice, "finish_reason", "stop") or "stop",
                usage=getattr(response, "usage", None),
            ),
            metadata={"provider": self.provider_name},
        )

    async def stream(self, request: ChatRequest) -> Sequence[ChatResponse]:
        messages = _build_litellm_messages(request)
        model = request.model or self.model
        kwargs = self._common_kwargs()
        if request.temperature is not None:
            kwargs["temperature"] = request.temperature
        if request.max_tokens is not None:
            kwargs["max_tokens"] = request.max_tokens
        kwargs["stream"] = True

        chunks: List[ChatResponse] = []
        accumulated = ""
        finish_reason = "stop"

        if self._router is not None:
            stream = self._router.acompletion(model=model, messages=messages, **kwargs)
        else:
            stream = self._litellm.acompletion(model=model, messages=messages, **kwargs)
        # LiteLLM returns the async generator directly when stream=True.
        if hasattr(stream, "__await__"):
            stream = await stream

        async for chunk in stream:
            delta = getattr(chunk.choices[0], "delta", None)
            piece = getattr(delta, "content", None) if delta is not None else None
            if isinstance(piece, str) and piece:
                accumulated += piece
            fr = getattr(chunk.choices[0], "finish_reason", None)
            if fr is not None:
                finish_reason = fr or "stop"
            chunks.append(
                ChatResponse(
                    choice=ChatChoice(
                        message=ChatMessage(
                            role="assistant",
                            content=accumulated,
                            metadata={"tenant": self.tenant, "provider": self.provider_name},
                        ),
                        finish_reason=finish_reason,
                    ),
                    metadata={"tenant": self.tenant, "streaming": True, "provider": self.provider_name},
                )
            )
        return tuple(chunks)

    async def models(self) -> Sequence[ModelInfo]:
        if self._router is not None:
            try:
                deployment_names = [d for d in self._router.get_model_names()]  # type: ignore[attr-defined]
            except Exception:  # pragma: no cover - router API variance
                deployment_names = [self.model]
            return [
                ModelInfo(id=name, provider=self.provider_name, metadata={"tenant": self.tenant})
                for name in deployment_names
            ]
        return [ModelInfo(id=self.model, provider=self.provider_name, metadata={"tenant": self.tenant})]

ModelInfo dataclass

Source code in src/ciel/providers/__init__.py
@dataclass(frozen=True)
class ModelInfo:
    id: str
    provider: str
    capabilities: Sequence[str] = ()
    context_window: Optional[int] = None
    metadata: Dict[str, Any] = field(default_factory=dict)

ProviderError

Bases: CielError

Model/provider error.

Source code in src/ciel/common/__init__.py
class ProviderError(CielError):
    """Model/provider error."""

_build_litellm_messages(request: ChatRequest) -> List[Dict[str, Any]]

Build LiteLLM's messages payload from a :class:ChatRequest.

Reuses the OpenAI-compatible content normalization so multimodal content (text + image_url parts) is forwarded unchanged.

Source code in src/ciel/providers/litellm.py
def _build_litellm_messages(request: ChatRequest) -> List[Dict[str, Any]]:
    """Build LiteLLM's ``messages`` payload from a :class:`ChatRequest`.

    Reuses the OpenAI-compatible content normalization so multimodal content
    (text + image_url parts) is forwarded unchanged.
    """
    from ciel.providers import _normalize_content_parts

    messages: List[Dict[str, Any]] = []
    for message in request.messages:
        payload: Dict[str, Any] = {
            "role": message.role,
            "content": _normalize_content_parts(message.content),
        }
        if message.name:
            payload["name"] = message.name
        messages.append(payload)
    return messages

_import_litellm()

Import litellm on demand; raise a clear error if it is missing.

Source code in src/ciel/providers/litellm.py
def _import_litellm():
    """Import ``litellm`` on demand; raise a clear error if it is missing."""
    try:
        import litellm  # type: ignore
    except ImportError as exc:  # pragma: no cover - depends on extras install
        raise ProviderError(
            "The 'litellm' extra is required for LiteLLMProvider. "
            "Install it with: pip install \"mana-ciel[litellm]\""
        ) from exc
    return litellm

Google Gemini provider (Generative Language API).

Offline-safe: only performs HTTP when an api_key is configured; in tests a mock httpx client can be injected. Mirrors the ChatProvider contract used by the OpenAI/Anthropic providers.

ChatChoice dataclass

Source code in src/ciel/runtime/tools.py
@dataclass(frozen=True)
class ChatChoice:
    message: ChatMessage
    finish_reason: str
    usage: Optional[Dict[str, Any]] = None
    metadata: Dict[str, Any] = field(default_factory=dict)

ChatMessage dataclass

Source code in src/ciel/runtime/tools.py
@dataclass(frozen=True)
class ChatMessage:
    role: str
    content: ChatContent
    name: Optional[str] = None
    tool_call_id: Optional[str] = None
    tool_calls: Optional[list[dict[str, Any]]] = None
    metadata: Dict[str, Any] = field(default_factory=dict)

    def text(self) -> str:
        """Extract plain text from content, tolerant to multimodal parts.

        - ``str`` content is returned verbatim.
        - ``list`` content concatenates the ``text`` of every part whose type
          is ``"text"`` but drops images/audio/video, so consumers (CLI,
          compression, ``AgentResponse.text``) see only readable text.
        """
        content = self.content
        if isinstance(content, str):
            return content
        if not isinstance(content, list):
            return ""
        parts: list[str] = []
        for part in content:
            if isinstance(part, dict) and part.get("type") == "text":
                text = part.get("text", "")
                if isinstance(text, str):
                    parts.append(text)
        return "".join(parts)

text() -> str

Extract plain text from content, tolerant to multimodal parts.

  • str content is returned verbatim.
  • list content concatenates the text of every part whose type is "text" but drops images/audio/video, so consumers (CLI, compression, AgentResponse.text) see only readable text.
Source code in src/ciel/runtime/tools.py
def text(self) -> str:
    """Extract plain text from content, tolerant to multimodal parts.

    - ``str`` content is returned verbatim.
    - ``list`` content concatenates the ``text`` of every part whose type
      is ``"text"`` but drops images/audio/video, so consumers (CLI,
      compression, ``AgentResponse.text``) see only readable text.
    """
    content = self.content
    if isinstance(content, str):
        return content
    if not isinstance(content, list):
        return ""
    parts: list[str] = []
    for part in content:
        if isinstance(part, dict) and part.get("type") == "text":
            text = part.get("text", "")
            if isinstance(text, str):
                parts.append(text)
    return "".join(parts)

ChatProvider

Bases: ABC

Source code in src/ciel/providers/__init__.py
class ChatProvider(ABC):
    provider_name: str = "unknown"

    @abstractmethod
    async def complete(self, request: "ChatRequest") -> "ChatResponse":
        raise NotImplementedError

    @abstractmethod
    async def stream(self, request: "ChatRequest") -> Sequence["ChatResponse"]:
        raise NotImplementedError

    @abstractmethod
    async def models(self) -> Sequence[ModelInfo]:
        raise NotImplementedError

ChatRequest dataclass

Source code in src/ciel/runtime/tools.py
@dataclass(frozen=True)
class ChatRequest:
    messages: Sequence[ChatMessage]
    tools: Sequence[ToolSpec] = ()
    model: Optional[str] = None
    temperature: Optional[float] = None
    max_tokens: Optional[int] = None
    extra: Dict[str, Any] = field(default_factory=dict)

ChatResponse dataclass

Source code in src/ciel/runtime/tools.py
@dataclass(frozen=True)
class ChatResponse:
    choice: ChatChoice
    metadata: Dict[str, Any] = field(default_factory=dict)

GeminiProvider

Bases: ChatProvider

Source code in src/ciel/providers/gemini.py
class GeminiProvider(ChatProvider):
    provider_name = "gemini"

    def __init__(
        self,
        *,
        api_key: Optional[str] = None,
        default_model: str = "gemini-1.5-flash",
        tenant: Optional[str] = None,
        base_url: str = "https://generativelanguage.googleapis.com/v1beta",
        timeout: float = 30.0,
        client: Optional[httpx.AsyncClient] = None,
    ) -> None:
        self.api_key = api_key
        self.default_model = default_model
        self.tenant = tenant
        self.base_url = base_url.rstrip("/")
        self.timeout = timeout
        self._client = client

    def _client_ctx(self):
        if self._client is not None:
            return _NullCtx(self._client)
        return httpx.AsyncClient(base_url=self.base_url, timeout=self.timeout)

    async def complete(self, request: ChatRequest) -> ChatResponse:
        if self.api_key is None:
            raise ProviderError("Gemini provider requires api_key")
        model = request.model or self.default_model
        payload = self._to_gemini(request)
        headers = {"Content-Type": "application/json", "x-goog-api-key": self.api_key}
        async with self._client_ctx() as client:
            response = await client.post(
                f"/models/{model}:generateContent", headers=headers, json=payload
            )
            response.raise_for_status()
            body = response.json()
        text = self._extract_text(body)
        message = ChatMessage(role="assistant", content=text, metadata={"tenant": self.tenant})
        return ChatResponse(
            choice=ChatChoice(message=message, finish_reason="stop"),
            metadata={"provider": self.provider_name},
        )

    async def stream(self, request: ChatRequest) -> Sequence[ChatResponse]:
        # Gemini streaming would use :streamGenerateContent; fall back to complete.
        return [await self.complete(request)]

    async def models(self) -> Sequence[ModelInfo]:
        return [ModelInfo(id=self.default_model, provider=self.provider_name, metadata={"tenant": self.tenant})]

    @staticmethod
    def _to_gemini(request: ChatRequest) -> dict[str, Any]:
        contents = []
        for message in request.messages:
            role = "model" if message.role == "assistant" else "user"
            contents.append({"role": role, "parts": _gemini_parts(message.content)})
        payload: dict[str, Any] = {"contents": contents}
        gen_cfg = {}
        if request.temperature is not None:
            gen_cfg["temperature"] = request.temperature
        if request.max_tokens is not None:
            gen_cfg["maxOutputTokens"] = request.max_tokens
        if gen_cfg:
            payload["generationConfig"] = gen_cfg
        return payload

    @staticmethod
    def _extract_text(body: dict[str, Any]) -> str:
        candidates = body.get("candidates") or []
        if not candidates:
            return ""
        parts = candidates[0].get("content", {}).get("parts", []) or []
        return "".join(part.get("text", "") for part in parts if isinstance(part, dict))

ModelInfo dataclass

Source code in src/ciel/providers/__init__.py
@dataclass(frozen=True)
class ModelInfo:
    id: str
    provider: str
    capabilities: Sequence[str] = ()
    context_window: Optional[int] = None
    metadata: Dict[str, Any] = field(default_factory=dict)

ProviderError

Bases: CielError

Model/provider error.

Source code in src/ciel/common/__init__.py
class ProviderError(CielError):
    """Model/provider error."""

_NullCtx

Context manager that yields an injected client without closing it.

Source code in src/ciel/providers/gemini.py
class _NullCtx:
    """Context manager that yields an injected client without closing it."""

    def __init__(self, client: httpx.AsyncClient) -> None:
        self.client = client

    async def __aenter__(self) -> httpx.AsyncClient:
        return self.client

    async def __aexit__(self, *exc: Any) -> None:
        return None

_gemini_parts(content: Any) -> List[Dict[str, Any]]

Normalize str | list[dict] content to Gemini parts.

  • str -> [{"text": <content>}]
  • list -> maps each part to a Gemini part:
    • {"type": "text", "text": ...} -> {"text": ...}
    • {"type": "image_url", "image_url": {"url": "data:..."}} -> {"inline_data": {"mime_type": ..., "data": <b64>}}
Source code in src/ciel/providers/gemini.py
def _gemini_parts(content: Any) -> List[Dict[str, Any]]:
    """Normalize ``str | list[dict]`` content to Gemini parts.

    - ``str`` -> ``[{"text": <content>}]``
    - ``list`` -> maps each part to a Gemini part:
        * ``{"type": "text", "text": ...}`` -> ``{"text": ...}``
        * ``{"type": "image_url", "image_url": {"url": "data:..."}}`` ->
          ``{"inline_data": {"mime_type": ..., "data": <b64>}}``
    """
    if isinstance(content, str):
        return [{"text": content}]
    if not isinstance(content, list):
        return [{"text": ""}]
    parts: List[Dict[str, Any]] = []
    for part in content:
        if not isinstance(part, dict):
            continue
        ptype = part.get("type")
        if ptype == "text":
            text = part.get("text", "")
            if isinstance(text, str):
                parts.append({"text": text})
        elif ptype == "image_url":
            url = (part.get("image_url") or {}).get("url", "")
            mime, data = _split_data_url(url)
            if mime is not None and data is not None:
                parts.append({"inline_data": {"mime_type": mime, "data": data}})
    return parts or [{"text": ""}]

_split_data_url(url: str) -> 'tuple[Optional[str], Optional[str]]'

Split a data:<mime>;base64,<data> URL into (mime, b64).

Source code in src/ciel/providers/gemini.py
def _split_data_url(url: str) -> "tuple[Optional[str], Optional[str]]":
    """Split a ``data:<mime>;base64,<data>`` URL into (mime, b64)."""
    if not isinstance(url, str) or not url.startswith("data:"):
        return None, None
    try:
        meta, data = url[5:].split(",", 1)
    except ValueError:
        return None, None
    if ";base64" not in meta:
        return None, None
    mime = meta.split(";base64", 1)[0] or "application/octet-stream"
    return mime, data

MockProvider (offline, sin extra)

Proveedor determinista para tests y evaluación sin red ni API keys. Tres modos:

  • fixed: responde una cadena constante (response=).
  • echo: repite la última palabra del prompt del usuario.
  • map: dict prompt -> response (coincidencia exacta o substring, case-insensitive).
from ciel.providers import MockProvider

# Respuesta fija (ideal para eval de métricas cerradas)
provider = MockProvider(mode="fixed", response="París")

# Mapa de prompts a respuestas
provider = MockProvider(mode="map", mapping={"capital de Francia": "París"})

# Auto-provider por prefijo (Fase 18)
from ciel.providers.auto import auto_provider
p = auto_provider("mock/echo")  # -> MockProvider(mode="echo")

Es útil como evaluable de ciel.evaluate (ver ciel.eval) para correr datasets de forma reproducible en CI sin consumir tokens.

MockProvider determinista para tests y evaluación offline (Fase 18).

Devuelve respuestas configurables sin red ni API keys, con tres modos:

  • echo: repite la última palabra del prompt del usuario (paridad con _EchoProvider de la CLI, que devuelve echo:<prompt> completo). Aquí repetimos la última palabra para que sea útil en eval de métricas cerradas.
  • map: dict prompt -> response (coincidencia exacta o substring).
  • fixed: respuesta constante.

Soporta complete / stream (parity: devuelve (response,)) / models y respeta el contrato ChatProvider (ciel.providers.ChatProvider). Offline-safe por construcción: no importa nada de red.

ChatChoice dataclass

Source code in src/ciel/runtime/tools.py
@dataclass(frozen=True)
class ChatChoice:
    message: ChatMessage
    finish_reason: str
    usage: Optional[Dict[str, Any]] = None
    metadata: Dict[str, Any] = field(default_factory=dict)

ChatMessage dataclass

Source code in src/ciel/runtime/tools.py
@dataclass(frozen=True)
class ChatMessage:
    role: str
    content: ChatContent
    name: Optional[str] = None
    tool_call_id: Optional[str] = None
    tool_calls: Optional[list[dict[str, Any]]] = None
    metadata: Dict[str, Any] = field(default_factory=dict)

    def text(self) -> str:
        """Extract plain text from content, tolerant to multimodal parts.

        - ``str`` content is returned verbatim.
        - ``list`` content concatenates the ``text`` of every part whose type
          is ``"text"`` but drops images/audio/video, so consumers (CLI,
          compression, ``AgentResponse.text``) see only readable text.
        """
        content = self.content
        if isinstance(content, str):
            return content
        if not isinstance(content, list):
            return ""
        parts: list[str] = []
        for part in content:
            if isinstance(part, dict) and part.get("type") == "text":
                text = part.get("text", "")
                if isinstance(text, str):
                    parts.append(text)
        return "".join(parts)

text() -> str

Extract plain text from content, tolerant to multimodal parts.

  • str content is returned verbatim.
  • list content concatenates the text of every part whose type is "text" but drops images/audio/video, so consumers (CLI, compression, AgentResponse.text) see only readable text.
Source code in src/ciel/runtime/tools.py
def text(self) -> str:
    """Extract plain text from content, tolerant to multimodal parts.

    - ``str`` content is returned verbatim.
    - ``list`` content concatenates the ``text`` of every part whose type
      is ``"text"`` but drops images/audio/video, so consumers (CLI,
      compression, ``AgentResponse.text``) see only readable text.
    """
    content = self.content
    if isinstance(content, str):
        return content
    if not isinstance(content, list):
        return ""
    parts: list[str] = []
    for part in content:
        if isinstance(part, dict) and part.get("type") == "text":
            text = part.get("text", "")
            if isinstance(text, str):
                parts.append(text)
    return "".join(parts)

ChatProvider

Bases: ABC

Source code in src/ciel/providers/__init__.py
class ChatProvider(ABC):
    provider_name: str = "unknown"

    @abstractmethod
    async def complete(self, request: "ChatRequest") -> "ChatResponse":
        raise NotImplementedError

    @abstractmethod
    async def stream(self, request: "ChatRequest") -> Sequence["ChatResponse"]:
        raise NotImplementedError

    @abstractmethod
    async def models(self) -> Sequence[ModelInfo]:
        raise NotImplementedError

ChatResponse dataclass

Source code in src/ciel/runtime/tools.py
@dataclass(frozen=True)
class ChatResponse:
    choice: ChatChoice
    metadata: Dict[str, Any] = field(default_factory=dict)

MockProvider

Bases: ChatProvider

Proveedor determinista configurable para eval/tests sin red.

Parameters:

Name Type Description Default
mode str

"echo" | "map" | "fixed".

'fixed'
response str

respuesta constante (modo fixed).

''
mapping Optional[Dict[str, str]]

dict prompt -> response (modo map). La coincidencia es exacta; si no hay exacta, se usa la primera clave que sea substring del prompt (case-insensitive).

None
model str

id de modelo devuelto por models() (default "mock").

'mock'
tenant Optional[str]

tenant_id propagado en metadata.

None
Source code in src/ciel/providers/mock.py
class MockProvider(ChatProvider):
    """Proveedor determinista configurable para eval/tests sin red.

    Args:
        mode: ``"echo"`` | ``"map"`` | ``"fixed"``.
        response: respuesta constante (modo ``fixed``).
        mapping: dict ``prompt -> response`` (modo ``map``). La coincidencia es
            exacta; si no hay exacta, se usa la primera clave que sea substring
            del prompt (case-insensitive).
        model: id de modelo devuelto por ``models()`` (default ``"mock"``).
        tenant: tenant_id propagado en metadata.
    """

    provider_name = "mock"

    def __init__(
        self,
        *,
        mode: str = "fixed",
        response: str = "",
        mapping: Optional[Dict[str, str]] = None,
        model: str = "mock",
        tenant: Optional[str] = None,
    ) -> None:
        if mode not in ("echo", "map", "fixed"):
            raise ValueError(f"MockProvider mode inválido: {mode!r}; usar echo|map|fixed")
        self.mode = mode
        self.response = response
        self.mapping = dict(mapping or {})
        self.model = model
        self.tenant = tenant

    def _respond(self, request: Any) -> str:
        if self.mode == "fixed":
            return self.response
        if self.mode == "echo":
            last = request.messages[-1].content if request.messages else ""
            if isinstance(last, list):
                last = "".join(p.get("text", "") for p in last if isinstance(p, dict))
            last = last or ""
            words = last.split()
            return words[-1] if words else ""
        # map
        prompt = request.messages[-1].content if request.messages else ""
        if isinstance(prompt, list):
            prompt = "".join(p.get("text", "") for p in prompt if isinstance(p, dict))
        prompt = prompt or ""
        if prompt in self.mapping:
            return self.mapping[prompt]
        lowered = prompt.lower()
        for key, val in self.mapping.items():
            if key and key.lower() in lowered:
                return val
        return self.mapping.get("", self.response)

    async def complete(self, request: Any) -> ChatResponse:
        text = self._respond(request)
        message = ChatMessage(role="assistant", content=text, metadata={"tenant": self.tenant})
        return ChatResponse(
            choice=ChatChoice(message=message, finish_reason="stop"),
            metadata={"provider": self.provider_name, "mode": self.mode, "tenant": self.tenant},
        )

    async def stream(self, request: Any) -> Sequence[ChatResponse]:
        return (await self.complete(request),)

    async def models(self) -> Sequence[ModelInfo]:
        return [
            ModelInfo(
                id=self.model,
                provider=self.provider_name,
                metadata={"tenant": self.tenant, "mode": self.mode},
            )
        ]

ModelInfo dataclass

Source code in src/ciel/providers/__init__.py
@dataclass(frozen=True)
class ModelInfo:
    id: str
    provider: str
    capabilities: Sequence[str] = ()
    context_window: Optional[int] = None
    metadata: Dict[str, Any] = field(default_factory=dict)