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. |
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
|
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)
|