Files

1382 lines
54 KiB
Python

"""Versioned provider adapter registry used by model configuration V3.
The registry is the only place where provider protocol differences are
described. Configuration may select an exact registered revision, but cannot
invent authentication headers, request paths, or provider-specific fields.
"""
from __future__ import annotations
import asyncio
import base64
import hashlib
import json
from collections.abc import Mapping
from dataclasses import asdict, dataclass, field
from typing import Any, Literal
from urllib.parse import quote
from .contracts import EvoRuntimeError
AdapterLifecycle = Literal["active", "deprecated", "blocked"]
_PROBE_RETRY_BASE_SECONDS = 0.25
_PROBE_RETRY_MAX_SECONDS = 5.0
_KIMI_CHAT_COMPLETION_TOKEN_MODELS = frozenset(
{
"k3",
"k3-256k",
"kimi-k3",
"kimi-for-coding",
"kimi-for-coding-highspeed",
}
)
@dataclass(frozen=True, slots=True)
class ErrorDisposition:
error_code: str
retryable: bool
health_scope: Literal["none", "provider_connection", "model_route"]
retry_after_ms: int | None = None
@dataclass(frozen=True, slots=True)
class NormalizedUsage:
input_tokens: int | None
output_tokens: int | None
cached_input_tokens: int | None
reasoning_tokens: int | None = None
total_tokens: int | None = None
provider_request_id_hash: str | None = None
finality: Literal["confirmed", "partial", "unconfirmed"] = "unconfirmed"
def confirmed_projection(self) -> Mapping[str, int | str] | None:
if self.finality != "confirmed" or any(
value is None
for value in (
self.input_tokens,
self.output_tokens,
self.cached_input_tokens,
)
):
return None
assert self.input_tokens is not None
assert self.output_tokens is not None
assert self.cached_input_tokens is not None
if (
min(self.input_tokens, self.output_tokens, self.cached_input_tokens) < 0
or self.cached_input_tokens > self.input_tokens
or (
self.reasoning_tokens is not None
and (
self.reasoning_tokens < 0
or self.reasoning_tokens > self.output_tokens
)
)
):
return None
result: dict[str, int | str] = {
"input_tokens": self.input_tokens,
"cached_input_tokens": self.cached_input_tokens,
"output_tokens": self.output_tokens,
}
if self.reasoning_tokens is not None:
result["reasoning_tokens"] = self.reasoning_tokens
if self.total_tokens is not None:
result["total_tokens"] = self.total_tokens
if self.provider_request_id_hash:
result["provider_request_id_hash"] = self.provider_request_id_hash
return result
@dataclass(frozen=True, slots=True)
class ParameterRule:
kind: Literal["boolean", "integer", "number", "string", "enum", "object"]
minimum: float | None = None
maximum: float | None = None
minimum_exclusive: bool = False
maximum_exclusive: bool = False
choices: tuple[str, ...] = ()
def validate(self, value: Any, path: str) -> None:
if self.kind == "boolean":
valid = isinstance(value, bool)
elif self.kind == "integer":
valid = isinstance(value, int) and not isinstance(value, bool)
elif self.kind == "number":
valid = isinstance(value, int | float) and not isinstance(value, bool)
elif self.kind in {"string", "enum"}:
valid = isinstance(value, str)
else:
valid = isinstance(value, Mapping)
if not valid:
raise EvoRuntimeError(
"MODEL_PARAMETER_INVALID",
f"{path} has an invalid type",
details=({"path": path, "code": "PARAMETER_TYPE_INVALID"},),
)
if self.kind == "enum" and value not in self.choices:
raise EvoRuntimeError(
"MODEL_PARAMETER_INVALID",
f"{path} is not supported",
details=({"path": path, "code": "PARAMETER_VALUE_UNSUPPORTED"},),
)
if self.kind in {"integer", "number"}:
number = float(value)
if self.minimum is not None and (
number < self.minimum
or (self.minimum_exclusive and number == self.minimum)
):
raise EvoRuntimeError(
"MODEL_PARAMETER_INVALID",
f"{path} is below its minimum",
details=({"path": path, "code": "PARAMETER_OUT_OF_RANGE"},),
)
if self.maximum is not None and (
number > self.maximum
or (self.maximum_exclusive and number == self.maximum)
):
raise EvoRuntimeError(
"MODEL_PARAMETER_INVALID",
f"{path} exceeds its maximum",
details=({"path": path, "code": "PARAMETER_OUT_OF_RANGE"},),
)
@dataclass(frozen=True, slots=True)
class ModelDescriptor:
context_tokens: int
max_output_tokens: int
capabilities: frozenset[str]
parameters: Mapping[str, ParameterRule]
@dataclass(frozen=True, slots=True)
class ModelDiscoveryDescriptor:
context_tokens: int | None
max_output_tokens: int | None
capabilities: frozenset[str]
reasoning_mode: Literal["none", "boolean", "effort"] = "none"
reasoning_efforts: tuple[str, ...] = ()
default_reasoning_effort: str | None = None
source: str = "adapter_catalog"
@dataclass(frozen=True, slots=True)
class AdapterRegistration:
adapter_id: str
adapter_revision: str
display_name: str
lifecycle: AdapterLifecycle
supported_wire_protocols: tuple[str, ...]
supported_api_modes: tuple[str, ...]
recommended_api_mode: str
recommended_base_url: str
auth_profile: Literal["bearer", "anthropic_api_key", "google_api_key"]
runtime_provider: str
protocol_hard_limits: Mapping[str, tuple[int, int]]
parameter_schema: Mapping[str, ParameterRule]
legacy_parameter_schema: Mapping[str, ParameterRule] = field(default_factory=dict)
reasoning_mode: Literal["none", "boolean", "effort"] = "none"
exact_model_overrides: Mapping[str, ModelDescriptor] = field(default_factory=dict)
discovery_model_overrides: Mapping[str, ModelDiscoveryDescriptor] = field(
default_factory=dict
)
discovery_capability: bool = False
server_tools_configurable: bool = False
replacement_revision: str | None = None
implementation_fingerprint: str = ""
def metadata(self) -> Mapping[str, Any]:
return {
"adapter_id": self.adapter_id,
"adapter_revision": self.adapter_revision,
"implementation_fingerprint": self.implementation_fingerprint,
"lifecycle": self.lifecycle,
"replacement_revision": self.replacement_revision,
"display_name": self.display_name,
"wire_protocols": list(self.supported_wire_protocols),
"recommended_base_url": self.recommended_base_url,
"credential_kind": "api_key",
"api_modes": [
{"id": value, "recommended": value == self.recommended_api_mode}
for value in self.supported_api_modes
],
"server_tools_configurable": self.server_tools_configurable,
"reasoning_mode": self.reasoning_mode,
"parameter_schema": {
key: asdict(rule) for key, rule in sorted(self.parameter_schema.items())
},
}
@property
def all_parameter_schema(self) -> Mapping[str, ParameterRule]:
return {**self.legacy_parameter_schema, **self.parameter_schema}
def resolve_model_descriptor(
self,
provider_model_id: str,
api_mode: str,
*,
context_tokens: int | None,
max_output_tokens: int | None,
declared_capabilities: Mapping[str, bool],
) -> ModelDescriptor:
if self.lifecycle == "blocked":
raise EvoRuntimeError("MODEL_ADAPTER_BLOCKED")
if api_mode not in self.supported_api_modes:
raise EvoRuntimeError("MODEL_API_MODE_UNSUPPORTED")
hard_context, hard_output = self.protocol_hard_limits[api_mode]
exact = self.exact_model_overrides.get(provider_model_id)
descriptor_context = exact.context_tokens if exact else hard_context
descriptor_output = exact.max_output_tokens if exact else hard_output
if context_tokens is None or max_output_tokens is None:
if exact is None:
raise EvoRuntimeError("TOKEN_BOUND_UNAVAILABLE")
context_tokens = context_tokens or descriptor_context
max_output_tokens = max_output_tokens or descriptor_output
if context_tokens > descriptor_context or max_output_tokens > descriptor_output:
raise EvoRuntimeError("MODEL_TOKEN_BOUND_EXCEEDED")
if max_output_tokens > context_tokens:
raise EvoRuntimeError("MODEL_TOKEN_BOUND_EXCEEDED")
# Exact descriptors are an advisory catalog: they supply known limits,
# defaults and UI hints, but do not overrule an administrator's product
# capability policy. A model catalog necessarily lags new model IDs and
# Provider releases; protocol validation happens when a real request is
# compiled instead.
supported = exact.capabilities if exact else frozenset({"text"})
return ModelDescriptor(
context_tokens=context_tokens,
max_output_tokens=max_output_tokens,
capabilities=supported,
parameters={**self.parameter_schema, **(exact.parameters if exact else {})},
)
def resolve_discovery_descriptor(
self, provider_model_id: str
) -> ModelDiscoveryDescriptor | None:
exact = self.exact_model_overrides.get(provider_model_id)
if exact is not None:
supports_reasoning = (
"thinking" in exact.capabilities and self.reasoning_mode != "none"
)
return ModelDiscoveryDescriptor(
context_tokens=exact.context_tokens,
max_output_tokens=exact.max_output_tokens,
capabilities=exact.capabilities,
reasoning_mode=self.reasoning_mode if supports_reasoning else "none",
reasoning_efforts=("low", "medium", "high")
if supports_reasoning
else (),
default_reasoning_effort="medium" if supports_reasoning else None,
source="adapter_exact_model",
)
return self.discovery_model_overrides.get(provider_model_id)
def validate_parameters(self, values: Mapping[str, Any], *, path: str) -> None:
schema = self.all_parameter_schema
unknown = set(values) - set(schema)
if unknown:
raise EvoRuntimeError(
"MODEL_PARAMETER_INVALID",
f"{path} contains unsupported parameters: {', '.join(sorted(unknown))}",
details=tuple(
{"path": f"{path}.{name}", "code": "PARAMETER_UNSUPPORTED"}
for name in sorted(unknown)
),
)
for name, value in values.items():
schema[name].validate(value, f"{path}.{name}")
if "temperature" in values and "top_p" in values:
raise EvoRuntimeError(
"MODEL_PARAMETER_CONFLICT",
f"{path} sets temperature and top_p",
details=(
{"path": f"{path}.temperature", "code": "PARAMETER_CONFLICT"},
{"path": f"{path}.top_p", "code": "PARAMETER_CONFLICT"},
),
)
reasoning = values.get("reasoning")
if reasoning is None and "reasoning_effort" in values:
reasoning = (
"off"
if values.get("reasoning_effort") == "disabled"
else values.get("reasoning_effort")
)
if reasoning is None and "thinking_enabled" in values:
reasoning = "high" if values.get("thinking_enabled") else "off"
if self.adapter_id == "dashscope":
forced_nonthinking = bool(values.get("structured_output")) or (
values.get("tool_choice") == "required"
)
if forced_nonthinking and reasoning not in {None, "off"}:
raise EvoRuntimeError(
"MODEL_PARAMETER_CONFLICT",
f"{path} combines reasoning with a Provider non-thinking mode",
details=(
{
"path": f"{path}.reasoning",
"code": "PARAMETER_CONFLICT",
},
),
)
if "reasoning_budget_tokens" in values and reasoning in {None, "off"}:
raise EvoRuntimeError(
"MODEL_PARAMETER_CONFLICT",
f"{path}.reasoning_budget_tokens requires reasoning",
details=(
{
"path": f"{path}.reasoning_budget_tokens",
"code": "PARAMETER_CONFLICT",
},
),
)
def compile_runtime_parameters(
self,
api_mode: str,
values: Mapping[str, Any],
output_token_limit: int,
*,
provider_model_id: str | None = None,
) -> Mapping[str, Any]:
"""Map canonical V3 parameters into installed LangChain constructors."""
result = dict(values)
result.pop("output_token_limit", None)
legacy_thinking = result.pop("thinking_enabled", None)
structured = result.pop("structured_output", None)
legacy_effort = result.pop("reasoning_effort", None)
reasoning = result.pop("reasoning", None)
reasoning_budget = result.pop("reasoning_budget_tokens", None)
if reasoning is None and legacy_effort is not None:
reasoning = "off" if legacy_effort == "disabled" else legacy_effort
if reasoning is None and legacy_thinking is not None:
reasoning = "high" if legacy_thinking else "off"
thinking = None if reasoning is None else reasoning != "off"
if self.adapter_id == "dashscope" and api_mode == "chat_completions":
forced_nonthinking = bool(structured) or result.get("tool_choice") == "required"
if forced_nonthinking and thinking:
raise EvoRuntimeError(
"MODEL_PARAMETER_CONFLICT",
"DashScope structured output and forced tools require thinking to be disabled",
)
if forced_nonthinking:
thinking = False
extra_body = {"enable_thinking": thinking} if thinking is not None else {}
if thinking and reasoning_budget is not None:
if int(reasoning_budget) > output_token_limit:
raise EvoRuntimeError(
"MODEL_PARAMETER_CONFLICT",
"reasoning budget exceeds the effective output limit",
)
extra_body["thinking_budget"] = int(reasoning_budget)
result.update({"max_completion_tokens": output_token_limit})
if extra_body:
result["extra_body"] = extra_body
elif self.adapter_id == "anthropic":
result["max_tokens"] = output_token_limit
if thinking:
if output_token_limit <= 1_024:
raise EvoRuntimeError(
"MODEL_PARAMETER_CONFLICT",
"Anthropic thinking requires max_tokens greater than 1024",
)
result["thinking"] = {
"type": "enabled",
"budget_tokens": int(reasoning_budget)
if reasoning_budget is not None
else min(8_192, max(1_024, output_token_limit // 2)),
}
elif api_mode == "responses":
result.update(
{
"max_output_tokens": output_token_limit,
"use_responses_api": True,
"store": False,
}
)
if reasoning not in {None, "off"}:
result["reasoning"] = {"effort": reasoning}
if self.adapter_id in {"openai", "xai"}:
result["reasoning"]["summary"] = "auto"
elif self.adapter_id == "dashscope" and thinking:
result["reasoning"] = {"effort": "medium"}
elif self.adapter_id == "google-gemini":
result["max_output_tokens"] = output_token_limit
if api_mode == "interactions":
result["store"] = False
if thinking is not None:
result["thinking"] = thinking
else:
if self._uses_max_completion_tokens(provider_model_id, api_mode):
result["max_completion_tokens"] = output_token_limit
else:
result["max_tokens"] = output_token_limit
# V3 routes declare the API envelope explicitly. Do not leave
# LangChain to infer the Responses API from model naming or kwargs.
if self.runtime_provider == "openai":
result["use_responses_api"] = False
if reasoning not in {None, "off"}:
result["reasoning_effort"] = reasoning
if structured:
result["response_format"] = {"type": "json_object"}
return result
def _uses_max_completion_tokens(
self, provider_model_id: str | None, api_mode: str
) -> bool:
"""Return whether an OpenAI Chat model rejects legacy ``max_tokens``."""
model_id = str(provider_model_id or "").strip().lower()
normalized = model_id.replace("-", "")
return (
self.adapter_id == "openai"
and api_mode == "chat_completions"
and (
normalized.startswith("gpt5")
or model_id in _KIMI_CHAT_COMPLETION_TOKEN_MODELS
)
)
def classify_error(self, error: BaseException) -> ErrorDisposition:
"""Classify typed SDK/HTTP failures without inspecting user-facing text."""
try:
import httpx
except ImportError: # pragma: no cover - httpx is a runtime dependency
httpx = None # type: ignore[assignment]
status = getattr(error, "status_code", None)
if status is None:
response = getattr(error, "response", None)
status = getattr(response, "status_code", None)
retry_after_ms = None
headers = getattr(getattr(error, "response", None), "headers", None)
if headers:
try:
retry_after_ms = (
int(float(headers.get("retry-after", 0)) * 1000) or None
)
except (TypeError, ValueError):
retry_after_ms = None
if status in {401, 403}:
return ErrorDisposition(
"MODEL_AUTHENTICATION_FAILED", False, "provider_connection"
)
if status == 404:
return ErrorDisposition("MODEL_NOT_FOUND", False, "model_route")
if status in {400, 422}:
return ErrorDisposition(
"MODEL_PROVIDER_REQUEST_REJECTED", False, "model_route"
)
if status == 429:
return ErrorDisposition(
"MODEL_RATE_LIMITED", True, "provider_connection", retry_after_ms
)
if isinstance(status, int) and status >= 500:
return ErrorDisposition(
"MODEL_PROVIDER_ERROR", True, "provider_connection", retry_after_ms
)
if isinstance(error, (TimeoutError, ConnectionError)) or (
httpx is not None and isinstance(error, httpx.TimeoutException)
):
return ErrorDisposition("MODEL_TIMEOUT", True, "provider_connection")
if httpx is not None and isinstance(error, httpx.TransportError):
return ErrorDisposition("MODEL_PROVIDER_ERROR", True, "provider_connection")
if isinstance(error, json.JSONDecodeError):
return ErrorDisposition(
"MODEL_PROVIDER_RESPONSE_INVALID", True, "provider_connection"
)
if isinstance(error, EvoRuntimeError):
return ErrorDisposition(error.code, False, "none")
return ErrorDisposition("MODEL_PROVIDER_ERROR", False, "model_route")
async def discover_models(
self,
*,
base_url: str,
api_key: str,
timeout_seconds: int = 10,
) -> tuple[Mapping[str, str], ...]:
if not self.discovery_capability:
raise EvoRuntimeError("MODEL_DISCOVERY_UNSUPPORTED")
import httpx
root = base_url.rstrip("/")
if self.adapter_id == "anthropic":
url = root + "/v1/models"
headers = {
"x-api-key": api_key,
"anthropic-version": "2023-06-01",
}
elif self.adapter_id == "google-gemini":
url = root + "/v1beta/models"
headers = {"x-goog-api-key": api_key}
else:
url = root + "/models"
headers = {"Authorization": f"Bearer {api_key}"}
try:
async with httpx.AsyncClient(
timeout=httpx.Timeout(timeout_seconds),
follow_redirects=False,
trust_env=False,
) as client:
response = await client.get(url, headers=headers)
response.raise_for_status()
payload = response.json()
except Exception as exc:
disposition = self.classify_error(exc)
raise EvoRuntimeError(disposition.error_code) from exc
items = payload.get("data") if isinstance(payload, Mapping) else None
if items is None and isinstance(payload, Mapping):
items = payload.get("models")
result: list[Mapping[str, str]] = []
for item in items or []:
if not isinstance(item, Mapping):
continue
model_id = str(item.get("id") or item.get("name") or "").strip()
model_id = model_id.removeprefix("models/")
if not model_id or len(model_id) > 512:
continue
result.append(
{
"provider_model_id": model_id,
"display_name": str(
item.get("display_name") or item.get("displayName") or model_id
)[:512],
}
)
if len(result) >= 500:
break
return tuple(result)
def build_probe_request(
self,
*,
base_url: str,
api_key: str,
provider_model_id: str,
api_mode: str,
probe_kind: str,
) -> tuple[str, Mapping[str, str], Mapping[str, Any]]:
"""Compile a minimal real Provider request for one capability probe."""
if probe_kind in {"video", "documents"}:
raise EvoRuntimeError("MODEL_CAPABILITY_PROBE_UNSUPPORTED")
root = base_url.rstrip("/")
prompt = (
'Return exactly {"ok":true}.'
if probe_kind == "structured_output"
else "Call the probe_ok tool now."
if probe_kind == "tools"
else "Reply with the single word OK."
)
png = (
"iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAQAAAC1HAwCAAAAC0lEQVR4nGNg"
"YAAAAAMAASsJTYQAAAAASUVORK5CYII="
)
if self.auth_profile == "anthropic_api_key":
headers = {
"x-api-key": api_key,
"anthropic-version": "2023-06-01",
"content-type": "application/json",
}
elif self.auth_profile == "google_api_key":
headers = {"x-goog-api-key": api_key, "content-type": "application/json"}
else:
headers = {
"Authorization": f"Bearer {api_key}",
"content-type": "application/json",
}
tool = {
"name": "probe_ok",
"description": "Return a probe acknowledgement.",
"parameters": {"type": "object", "properties": {}},
}
if self.adapter_id == "anthropic":
content: Any = prompt
if probe_kind == "vision":
content = [
{
"type": "image",
"source": {
"type": "base64",
"media_type": "image/png",
"data": png,
},
},
{"type": "text", "text": prompt},
]
body: dict[str, Any] = {
"model": provider_model_id,
"messages": [{"role": "user", "content": content}],
"max_tokens": 16,
}
if probe_kind == "tools":
body["tools"] = [{**tool, "input_schema": tool["parameters"]}]
body["tools"][0].pop("parameters")
body["tool_choice"] = {"type": "tool", "name": "probe_ok"}
elif probe_kind == "structured_output":
body["output_config"] = {
"format": {
"type": "json_schema",
"schema": {
"type": "object",
"properties": {"ok": {"type": "boolean"}},
"required": ["ok"],
"additionalProperties": False,
},
}
}
elif probe_kind == "reasoning":
body["max_tokens"] = 1_025
body["thinking"] = {"type": "enabled", "budget_tokens": 1_024}
return root + "/v1/messages", headers, body
if self.adapter_id == "google-gemini":
if api_mode == "interactions":
interaction_content: list[dict[str, Any]] = [
{"type": "text", "text": prompt}
]
if probe_kind == "vision":
interaction_content.insert(
0,
{
"type": "image",
"mime_type": "image/png",
"data": png,
},
)
interaction_body: dict[str, Any] = {
"model": provider_model_id,
"input": [{"role": "user", "content": interaction_content}],
"generation_config": {"max_output_tokens": 16},
"store": False,
"stream": False,
}
if probe_kind == "tools":
interaction_body["tools"] = [{"type": "function", **tool}]
interaction_body["tool_choice"] = {
"type": "function",
"name": "probe_ok",
}
elif probe_kind == "structured_output":
interaction_body["generation_config"].update(
{
"response_mime_type": "application/json",
"response_schema": {
"type": "object",
"properties": {"ok": {"type": "boolean"}},
},
}
)
elif probe_kind == "reasoning":
interaction_body["generation_config"].update(
{"thinking_level": "high", "thinking_summaries": "auto"}
)
return root + "/v1beta/interactions", headers, interaction_body
parts: list[dict[str, Any]] = [{"text": prompt}]
if probe_kind == "vision":
parts.insert(
0,
{"inline_data": {"mime_type": "image/png", "data": png}},
)
body = {
"contents": [{"role": "user", "parts": parts}],
"generationConfig": {"maxOutputTokens": 16},
}
if probe_kind == "tools":
body["tools"] = [{"functionDeclarations": [tool]}]
body["toolConfig"] = {
"functionCallingConfig": {
"mode": "ANY",
"allowedFunctionNames": ["probe_ok"],
}
}
elif probe_kind == "structured_output":
body["generationConfig"].update(
{
"responseMimeType": "application/json",
"responseSchema": {
"type": "OBJECT",
"properties": {"ok": {"type": "BOOLEAN"}},
},
}
)
elif probe_kind == "reasoning":
body["generationConfig"]["thinkingConfig"] = {
"thinkingBudget": 1_024
}
return (
root
+ "/v1beta/models/"
+ quote(provider_model_id, safe="")
+ ":generateContent",
headers,
body,
)
if api_mode == "responses":
input_value: Any = prompt
if probe_kind == "vision":
input_value = [
{
"role": "user",
"content": [
{"type": "input_text", "text": prompt},
{
"type": "input_image",
"image_url": f"data:image/png;base64,{png}",
},
],
}
]
body = {
"model": provider_model_id,
"input": input_value,
"max_output_tokens": 16,
"store": False,
}
if probe_kind == "tools":
body["tools"] = [{"type": "function", **tool}]
body["tool_choice"] = {"type": "function", "name": "probe_ok"}
elif probe_kind == "structured_output":
body["text"] = {
"format": {
"type": "json_schema",
"name": "probe",
"schema": {
"type": "object",
"properties": {"ok": {"type": "boolean"}},
"required": ["ok"],
"additionalProperties": False,
},
"strict": True,
}
}
elif probe_kind == "reasoning":
body["reasoning"] = {"effort": "low"}
return root + "/responses", headers, body
content = prompt
if probe_kind == "vision":
content = [
{"type": "text", "text": prompt},
{
"type": "image_url",
"image_url": {"url": f"data:image/png;base64,{png}"},
},
]
body = {
"model": provider_model_id,
"messages": [{"role": "user", "content": content}],
}
if self.adapter_id == "dashscope" or self._uses_max_completion_tokens(
provider_model_id, api_mode
):
body["max_completion_tokens"] = 16
if self.adapter_id == "dashscope":
body["enable_thinking"] = False
else:
body["max_tokens"] = 16
if probe_kind == "tools":
body["tools"] = [{"type": "function", "function": tool}]
body["tool_choice"] = {
"type": "function",
"function": {"name": "probe_ok"},
}
elif probe_kind == "structured_output":
body["response_format"] = {"type": "json_object"}
elif probe_kind == "reasoning":
if self.adapter_id == "dashscope":
body["enable_thinking"] = True
body["thinking_budget"] = 16
body["stream"] = True
else:
body["reasoning_effort"] = "low"
return root + "/chat/completions", headers, body
async def probe_model(
self,
*,
base_url: str,
api_key: str,
provider_model_id: str,
api_mode: str,
probe_kinds: tuple[str, ...],
timeout_seconds: int,
max_attempts: int = 2,
) -> Mapping[str, str]:
import httpx
attempts_limit = max(1, min(int(max_attempts), 3))
results: dict[str, str] = {}
async with httpx.AsyncClient(
timeout=httpx.Timeout(timeout_seconds),
follow_redirects=False,
trust_env=False,
) as client:
for probe_kind in probe_kinds:
for attempt in range(1, attempts_limit + 1):
try:
url, headers, body = self.build_probe_request(
base_url=base_url,
api_key=api_key,
provider_model_id=provider_model_id,
api_mode=api_mode,
probe_kind=probe_kind,
)
async with client.stream(
"POST", url, headers=dict(headers), json=dict(body)
) as response:
response.raise_for_status()
chunks: list[bytes] = []
total = 0
async for chunk in response.aiter_bytes():
total += len(chunk)
if total > 1_048_576:
raise EvoRuntimeError(
"MODEL_PROVIDER_RESPONSE_INVALID"
)
chunks.append(chunk)
payloads = self._decode_probe_payloads(b"".join(chunks))
observed_revision = self._validate_probe_payloads(
probe_kind,
payloads,
)
except Exception as exc:
disposition = self.classify_error(exc)
if disposition.retryable and attempt < attempts_limit:
server_delay = (
disposition.retry_after_ms / 1_000
if disposition.retry_after_ms is not None
else _PROBE_RETRY_BASE_SECONDS * (2 ** (attempt - 1))
)
await asyncio.sleep(
min(
_PROBE_RETRY_MAX_SECONDS,
max(0.0, server_delay),
)
)
continue
raise EvoRuntimeError(
disposition.error_code,
details=(
{
"path": f"probe.{probe_kind}",
"code": disposition.error_code,
"probe_kind": probe_kind,
"attempts": attempt,
"retryable": disposition.retryable,
},
),
) from exc
results[probe_kind] = "supported"
if observed_revision:
results["resolved_model_revision"] = observed_revision
break
return results
async def test_model_connection(
self,
*,
base_url: str,
api_key: str,
provider_model_id: str,
api_mode: str,
timeout_seconds: int,
) -> Mapping[str, str]:
"""Make one minimal text request for an administrator connection test.
This intentionally does not infer or certify tools, reasoning, media,
or structured output. Those behaviours are exercised by real requests
and Adapter integration tests, never used as a model availability gate.
"""
return await self.probe_model(
base_url=base_url,
api_key=api_key,
provider_model_id=provider_model_id,
api_mode=api_mode,
probe_kinds=("connectivity",),
timeout_seconds=timeout_seconds,
)
@staticmethod
def _decode_probe_payloads(raw: bytes) -> tuple[Mapping[str, Any], ...]:
if not raw or len(raw) > 1_048_576:
raise EvoRuntimeError("MODEL_PROVIDER_RESPONSE_INVALID")
try:
text = raw.decode("utf-8")
except UnicodeDecodeError as exc:
raise EvoRuntimeError("MODEL_PROVIDER_RESPONSE_INVALID") from exc
payloads: list[Mapping[str, Any]] = []
stripped = text.strip()
try:
decoded = json.loads(stripped)
if isinstance(decoded, Mapping):
payloads.append(decoded)
except json.JSONDecodeError:
for line in text.splitlines():
if not line.startswith("data:"):
continue
data = line[5:].strip()
if not data or data == "[DONE]":
continue
try:
decoded = json.loads(data)
except json.JSONDecodeError as exc:
raise EvoRuntimeError("MODEL_PROVIDER_RESPONSE_INVALID") from exc
if isinstance(decoded, Mapping):
payloads.append(decoded)
if not payloads:
raise EvoRuntimeError("MODEL_PROVIDER_RESPONSE_INVALID")
return tuple(payloads)
@classmethod
def _validate_probe_payloads(
cls,
probe_kind: str,
payloads: tuple[Mapping[str, Any], ...],
) -> str:
if any(cls._contains_probe_key(payload, "error") for payload in payloads):
raise EvoRuntimeError("MODEL_PROVIDER_RESPONSE_INVALID")
recognized = any(
any(
key in payload
for key in (
"choices",
"content",
"candidates",
"output",
"outputs",
"response",
"id",
)
)
for payload in payloads
)
if not recognized:
raise EvoRuntimeError("MODEL_PROVIDER_RESPONSE_INVALID")
if probe_kind == "tools" and not any(
cls._contains_probe_key(payload, key)
for payload in payloads
for key in ("tool_calls", "function_call", "functionCall", "tool_use")
) and not any(
cls._contains_probe_value(payload, value)
for payload in payloads
for value in ("function_call", "tool_use")
):
raise EvoRuntimeError("MODEL_CAPABILITY_PROBE_FAILED")
if probe_kind == "structured_output":
structured = False
for payload in payloads:
for value in cls._probe_text_values(payload):
candidate = (
value.strip()
.removeprefix("```json")
.removesuffix("```")
.strip()
)
try:
parsed = json.loads(candidate)
except json.JSONDecodeError:
continue
if isinstance(parsed, Mapping) and isinstance(parsed.get("ok"), bool):
structured = True
break
if not structured:
raise EvoRuntimeError("MODEL_CAPABILITY_PROBE_FAILED")
if probe_kind == "reasoning" and not any(
cls._contains_probe_key(payload, key)
for payload in payloads
for key in (
"reasoning",
"reasoning_content",
"reasoning_tokens",
"thinking",
"thought",
"thoughts_token_count",
)
) and not any(
cls._contains_probe_value(payload, value)
for payload in payloads
for value in ("reasoning", "thinking")
):
raise EvoRuntimeError("MODEL_CAPABILITY_PROBE_FAILED")
for payload in reversed(payloads):
for key in (
"resolved_model_revision",
"modelVersion",
"model_version",
"model",
):
value = payload.get(key)
if isinstance(value, str) and value.strip():
return value.strip()[:512]
return ""
@classmethod
def _contains_probe_key(cls, value: Any, name: str) -> bool:
if isinstance(value, Mapping):
if name in value and value[name] not in (None, "", [], {}):
return True
return any(cls._contains_probe_key(item, name) for item in value.values())
if isinstance(value, list | tuple):
return any(cls._contains_probe_key(item, name) for item in value)
return False
@classmethod
def _contains_probe_value(cls, value: Any, expected: str) -> bool:
if isinstance(value, Mapping):
return any(
(isinstance(item, str) and item == expected)
or cls._contains_probe_value(item, expected)
for item in value.values()
)
if isinstance(value, list | tuple):
return any(cls._contains_probe_value(item, expected) for item in value)
return False
@classmethod
def _probe_text_values(cls, value: Any) -> tuple[str, ...]:
result: list[str] = []
if isinstance(value, Mapping):
for key, item in value.items():
if key in {"text", "content", "output_text"} and isinstance(item, str):
result.append(item)
else:
result.extend(cls._probe_text_values(item))
elif isinstance(value, list | tuple):
for item in value:
result.extend(cls._probe_text_values(item))
return tuple(result)
class AdapterRegistry:
def __init__(
self,
registrations: tuple[AdapterRegistration, ...],
*,
policy_key_id: str,
manifest_public_key: str,
manifest_signature: str,
) -> None:
self._registrations = {
(item.adapter_id, item.adapter_revision): item for item in registrations
}
if len(self._registrations) != len(registrations):
raise RuntimeError("duplicate adapter registration")
projection = [item.metadata() for item in registrations]
manifest = json.dumps(
projection, sort_keys=True, separators=(",", ":")
).encode()
try:
from cryptography.hazmat.primitives.asymmetric.ed25519 import (
Ed25519PublicKey,
)
if not manifest_signature:
raise ValueError("missing registry signature")
Ed25519PublicKey.from_public_bytes(
base64.b64decode(manifest_public_key, validate=True)
).verify(base64.b64decode(manifest_signature, validate=True), manifest)
except Exception as exc:
raise RuntimeError(
"adapter registry manifest signature is invalid"
) from exc
digest = hashlib.sha256(manifest).hexdigest()
self.registry_revision = f"sha256:{digest}"
self.policy_key_id = policy_key_id
self.manifest_signature = manifest_signature
def get(self, adapter_id: str, adapter_revision: str) -> AdapterRegistration:
item = self._registrations.get((adapter_id, adapter_revision))
if item is None:
raise EvoRuntimeError("MODEL_ADAPTER_UNAVAILABLE")
return item
def metadata(self) -> Mapping[str, Any]:
return {
"registry_revision": self.registry_revision,
"registry_policy_key_id": self.policy_key_id,
"manifest_signature": self.manifest_signature,
"adapters": [
item.metadata()
for item in sorted(
self._registrations.values(),
key=lambda value: (value.adapter_id, value.adapter_revision),
)
],
}
_LEGACY_REASONING = {
"reasoning_effort": ParameterRule(
"enum", choices=("disabled", "low", "medium", "high", "max")
),
"thinking_enabled": ParameterRule("boolean"),
}
_ANTHROPIC_PARAMETERS = {
"output_token_limit": ParameterRule("integer", minimum=1),
"temperature": ParameterRule(
"number", minimum=0, maximum=1, maximum_exclusive=False
),
"top_p": ParameterRule("number", minimum=0, maximum=1, minimum_exclusive=True),
"top_k": ParameterRule("integer", minimum=0),
"reasoning": ParameterRule(
"enum", choices=("off", "low", "medium", "high")
),
"reasoning_budget_tokens": ParameterRule("integer", minimum=1024),
"structured_output": ParameterRule("boolean"),
"tool_choice": ParameterRule("enum", choices=("auto", "none", "required")),
}
_OPENAI_PARAMETERS = {
"output_token_limit": ParameterRule("integer", minimum=1),
"temperature": ParameterRule(
"number", minimum=0, maximum=2, maximum_exclusive=True
),
"top_p": ParameterRule("number", minimum=0, maximum=1, minimum_exclusive=True),
"reasoning": ParameterRule(
"enum", choices=("off", "low", "medium", "high", "max")
),
"structured_output": ParameterRule("boolean"),
"tool_choice": ParameterRule("enum", choices=("auto", "none", "required")),
}
_GEMINI_PARAMETERS = {
"output_token_limit": ParameterRule("integer", minimum=1),
"temperature": ParameterRule(
"number", minimum=0, maximum=2, maximum_exclusive=True
),
"top_p": ParameterRule("number", minimum=0, maximum=1, minimum_exclusive=True),
"top_k": ParameterRule("integer", minimum=1),
"reasoning": ParameterRule(
"enum", choices=("off", "low", "medium", "high")
),
"structured_output": ParameterRule("boolean"),
"tool_choice": ParameterRule("enum", choices=("auto", "none", "required")),
}
_XAI_PARAMETERS = {
"output_token_limit": ParameterRule("integer", minimum=1),
"temperature": ParameterRule(
"number", minimum=0, maximum=2, maximum_exclusive=True
),
"top_p": ParameterRule("number", minimum=0, maximum=1, minimum_exclusive=True),
"reasoning": ParameterRule(
"enum", choices=("off", "low", "medium", "high")
),
"structured_output": ParameterRule("boolean"),
"tool_choice": ParameterRule("enum", choices=("auto", "none", "required")),
}
_DASHSCOPE_PARAMETERS = {
"output_token_limit": ParameterRule("integer", minimum=1, maximum=65_536),
"temperature": ParameterRule(
"number", minimum=0, maximum=2, maximum_exclusive=True
),
"top_p": ParameterRule("number", minimum=0, maximum=1, minimum_exclusive=True),
"reasoning": ParameterRule(
"enum", choices=("off", "low", "medium", "high")
),
"reasoning_budget_tokens": ParameterRule(
"integer", minimum=1, maximum=262_144
),
"structured_output": ParameterRule("boolean"),
"tool_choice": ParameterRule("enum", choices=("auto", "none", "required")),
}
_GENERIC_OPENAI_PARAMETERS = {
"output_token_limit": ParameterRule("integer", minimum=1),
"temperature": ParameterRule(
"number", minimum=0, maximum=2, maximum_exclusive=True
),
"top_p": ParameterRule("number", minimum=0, maximum=1, minimum_exclusive=True),
"structured_output": ParameterRule("boolean"),
"tool_choice": ParameterRule("enum", choices=("auto", "none", "required")),
}
def _fingerprint(adapter_id: str, revision: str) -> str:
return (
"sha256:"
+ hashlib.sha256(
f"{adapter_id}:{revision}:evoscientist-adapter-v1".encode()
).hexdigest()
)
def _registration(
adapter_id: str,
revision: str,
display_name: str,
wires: tuple[str, ...],
modes: tuple[str, ...],
recommended_mode: str,
base_url: str,
auth: Literal["bearer", "anthropic_api_key", "google_api_key"],
runtime_provider: str,
parameter_schema: Mapping[str, ParameterRule],
*,
reasoning_mode: Literal["none", "boolean", "effort"] = "none",
discovery: bool = True,
exact: Mapping[str, ModelDescriptor] | None = None,
discovery_models: Mapping[str, ModelDiscoveryDescriptor] | None = None,
) -> AdapterRegistration:
return AdapterRegistration(
adapter_id=adapter_id,
adapter_revision=revision,
display_name=display_name,
lifecycle="active",
supported_wire_protocols=wires,
supported_api_modes=modes,
recommended_api_mode=recommended_mode,
recommended_base_url=base_url,
auth_profile=auth,
runtime_provider=runtime_provider,
protocol_hard_limits=dict.fromkeys(modes, (2_000_000, 131_072)),
parameter_schema=parameter_schema,
legacy_parameter_schema=(
_LEGACY_REASONING if reasoning_mode != "none" else {}
),
reasoning_mode=reasoning_mode,
exact_model_overrides=exact or {},
discovery_model_overrides=discovery_models or {},
discovery_capability=discovery,
implementation_fingerprint=_fingerprint(adapter_id, revision),
)
_QWEN37 = ModelDescriptor(
context_tokens=1_000_000,
max_output_tokens=65_536,
# Advisory catalog defaults. Product policy still decides whether the
# application has implemented each input/output workflow.
capabilities=frozenset(
{"text", "vision", "tools", "thinking", "structured_output"}
),
parameters=_DASHSCOPE_PARAMETERS,
)
_KIMI_CODE_DISCOVERY_MODELS = {
"k3": ModelDiscoveryDescriptor(
context_tokens=1_048_576,
max_output_tokens=None,
capabilities=frozenset({"text", "thinking"}),
reasoning_mode="effort",
reasoning_efforts=("low", "high", "max"),
default_reasoning_effort="high",
source="kimi_code_official_catalog",
),
"kimi-for-coding": ModelDiscoveryDescriptor(
context_tokens=262_144,
max_output_tokens=None,
capabilities=frozenset({"text", "thinking"}),
reasoning_mode="effort",
reasoning_efforts=("high",),
default_reasoning_effort="high",
source="kimi_code_official_catalog",
),
"kimi-for-coding-highspeed": ModelDiscoveryDescriptor(
context_tokens=262_144,
max_output_tokens=None,
capabilities=frozenset({"text", "thinking"}),
reasoning_mode="effort",
reasoning_efforts=("high",),
default_reasoning_effort="high",
source="kimi_code_official_catalog",
),
}
BUILTIN_ADAPTER_REGISTRY = AdapterRegistry(
(
_registration(
"anthropic",
"anthropic-v1",
"Anthropic",
("anthropic_native",),
("messages",),
"messages",
"https://api.anthropic.com",
"anthropic_api_key",
"anthropic",
_ANTHROPIC_PARAMETERS,
reasoning_mode="boolean",
),
_registration(
"openai",
"openai-v1",
"OpenAI",
("openai_native",),
("responses", "chat_completions"),
"responses",
"https://api.openai.com/v1",
"bearer",
"openai",
_OPENAI_PARAMETERS,
reasoning_mode="effort",
discovery_models=_KIMI_CODE_DISCOVERY_MODELS,
),
_registration(
"google-gemini",
"google-gemini-v1",
"Google Gemini Developer API",
("gemini_native",),
("interactions", "generate_content"),
"interactions",
"https://generativelanguage.googleapis.com",
"google_api_key",
"google_genai",
_GEMINI_PARAMETERS,
reasoning_mode="effort",
),
_registration(
"xai",
"xai-v1",
"xAI Grok",
("openai_compatible",),
("responses", "chat_completions"),
"responses",
"https://api.x.ai/v1",
"bearer",
"openai",
_XAI_PARAMETERS,
reasoning_mode="effort",
),
_registration(
"dashscope",
"dashscope-v1",
"Alibaba Cloud DashScope",
("openai_compatible",),
("chat_completions", "responses"),
"chat_completions",
"https://dashscope.aliyuncs.com/compatible-mode/v1",
"bearer",
"openai",
_DASHSCOPE_PARAMETERS,
reasoning_mode="boolean",
exact={"qwen3.7-plus": _QWEN37, "qwen3.7-plus-2026-05-26": _QWEN37},
),
_registration(
"generic-openai-compatible",
"generic-openai-compatible-v1",
"Generic OpenAI Compatible",
("openai_compatible",),
("chat_completions",),
"chat_completions",
"https://example.invalid/v1",
"bearer",
"openai",
_GENERIC_OPENAI_PARAMETERS,
),
),
policy_key_id="evoscientist-adapter-registry-release-v2",
manifest_public_key="9NVTpwijNh5L4+yykcr6uoJukNGgihp8V8DGSHSXunA=",
manifest_signature="7/3X6ptDOf5KLxjEepI2jDsHPIu2l/2XfHXzRbveQVIPbD9cLp82OIUM9QEPejKzzWcQqRnNekNyOUy2OrLWCA==",
)
def get_adapter_registry() -> AdapterRegistry:
return BUILTIN_ADAPTER_REGISTRY