3889 lines
162 KiB
Python
3889 lines
162 KiB
Python
"""Transactional V3 model runtime owned by EvoScientist."""
|
|
|
|
from __future__ import annotations
|
|
|
|
import asyncio
|
|
import base64
|
|
import binascii
|
|
import hashlib
|
|
import json
|
|
import logging
|
|
import math
|
|
import re
|
|
import time
|
|
import traceback
|
|
import uuid
|
|
from collections import defaultdict
|
|
from collections.abc import AsyncIterator, Callable, Mapping, Sequence
|
|
from contextvars import ContextVar
|
|
from dataclasses import asdict, dataclass, field, replace
|
|
from typing import Any
|
|
|
|
from langchain_core.callbacks import AsyncCallbackHandler
|
|
from langchain_core.messages import BaseMessage, message_to_dict
|
|
|
|
from .adapter_registry import AdapterRegistration, NormalizedUsage, get_adapter_registry
|
|
from .configuration import SecretResolver
|
|
from .contracts import (
|
|
AdmissionGrant,
|
|
AdmissionGrantVerifier,
|
|
AgentInputV3,
|
|
AgentModelSet,
|
|
EvoRuntimeError,
|
|
EvoRuntimeEvent,
|
|
EvoWebRun,
|
|
HmacGrantAuthority,
|
|
ModelAttemptEvent,
|
|
ModelCatalog,
|
|
ModelCatalogEntry,
|
|
ModelFactoryResult,
|
|
PreparedRunQuote,
|
|
PricingQuote,
|
|
RouteCallBound,
|
|
RouteIdentity,
|
|
RoutePreparationGrant,
|
|
VerifiedModelSubject,
|
|
WebHostContext,
|
|
event_payload,
|
|
now_ms,
|
|
)
|
|
from .crypto import HmacKeyRing, canonical_json_v1
|
|
from .host_execution_registry import HostExecutionRegistry
|
|
from .invocation import InvocationPlan, compile_invocation_plan
|
|
from .model_config import (
|
|
EvoModelConfig,
|
|
FileEvoModelConfigStore,
|
|
RouteRef,
|
|
endpoint_fingerprint,
|
|
invocation_fingerprint,
|
|
resolve_secret,
|
|
route_fingerprint,
|
|
route_semantics_hash,
|
|
)
|
|
from .user_options import (
|
|
model_options_schema_hash,
|
|
project_user_options_for_purpose,
|
|
validate_user_model_options,
|
|
)
|
|
|
|
_BIGINT_MAX = 2**63 - 1
|
|
logger = logging.getLogger(__name__)
|
|
_construction_owner: ContextVar[Any] = ContextVar("evo_construction_owner", default=None)
|
|
_ROUTE_SEMANTICS_INFO = "ai4sci/route-semantics-hash/v3"
|
|
_ROUTE_FINGERPRINT_INFO = "ai4sci/route-fingerprint/v3"
|
|
_ENDPOINT_FINGERPRINT_INFO = "ai4sci/endpoint-fingerprint/v3"
|
|
_PROTOCOL_MARGIN_TOKENS: Mapping[tuple[str, str], int] = {
|
|
("custom-openai", "chat_completions"): 64,
|
|
("custom-openai", "responses"): 96,
|
|
# Generic OpenAI-compatible providers use the Chat Completions message
|
|
# envelope and therefore have the same conservative framing margin.
|
|
("generic-openai-compatible", "chat_completions"): 64,
|
|
("openai", "chat_completions"): 64,
|
|
("openai", "responses"): 96,
|
|
("anthropic", "messages"): 64,
|
|
("dashscope", "chat_completions"): 64,
|
|
("dashscope", "responses"): 96,
|
|
("google-gemini", "interactions"): 96,
|
|
("google-gemini", "generate_content"): 64,
|
|
("xai", "responses"): 96,
|
|
("xai", "chat_completions"): 64,
|
|
}
|
|
|
|
_DEBUG_INVOCATION_PARAMETER_KEYS = frozenset(
|
|
{
|
|
"disable_streaming",
|
|
"extra_body",
|
|
"max_completion_tokens",
|
|
"max_output_tokens",
|
|
"max_tokens",
|
|
"reasoning",
|
|
"reasoning_effort",
|
|
"response_format",
|
|
"store",
|
|
"streaming",
|
|
"temperature",
|
|
"thinking",
|
|
"tool_choice",
|
|
"top_p",
|
|
"use_responses_api",
|
|
}
|
|
)
|
|
_DEBUG_NESTED_PARAMETER_KEYS = frozenset(
|
|
{
|
|
"budget_tokens",
|
|
"effort",
|
|
"enable_thinking",
|
|
"thinking_budget",
|
|
"type",
|
|
}
|
|
)
|
|
_AUTH_VALUE_PATTERN = re.compile(
|
|
r"(?i)(bearer\s+|(?:api[_-]?key|authorization|token|secret)\s*[:=]\s*)"
|
|
r"([^\s,;}'\"]+)"
|
|
)
|
|
_OPENAI_RESPONSE_EVENT_TYPE_PATTERN = re.compile(
|
|
r'"type"\s*:\s*"(response\.[A-Za-z0-9_.:-]{1,120})"'
|
|
)
|
|
|
|
|
|
def _enabled_purposes(title_policy: str) -> tuple[str, ...]:
|
|
purposes = ["main_agent", "tool_selector", "deepagents_summarizer"]
|
|
if title_policy == "best_effort":
|
|
purposes.append("title")
|
|
return tuple(purposes)
|
|
|
|
|
|
def _merge_params(
|
|
base: Mapping[str, Any], overlay: Mapping[str, Any]
|
|
) -> dict[str, Any]:
|
|
"""Merge provider adapter parameters without discarding sibling JSON fields."""
|
|
|
|
result = dict(base)
|
|
for key, value in overlay.items():
|
|
existing = result.get(key)
|
|
if isinstance(existing, Mapping) and isinstance(value, Mapping):
|
|
result[key] = _merge_params(existing, value)
|
|
else:
|
|
result[key] = value
|
|
return result
|
|
|
|
|
|
@dataclass(frozen=True, slots=True)
|
|
class ResolvedRoute:
|
|
ref: RouteRef
|
|
identity: RouteIdentity
|
|
quote: PricingQuote
|
|
context_window: int
|
|
max_output_tokens: int
|
|
reasoning_mode: str
|
|
reasoning_enabled_params: Mapping[str, Any]
|
|
reasoning_disabled_params: Mapping[str, Any]
|
|
base_url: str
|
|
params: Mapping[str, Any]
|
|
api_key: str
|
|
default_headers: Mapping[str, str]
|
|
secret_fingerprints: Mapping[str, str]
|
|
adapter: AdapterRegistration | None = None
|
|
provider_max_inflight: int = 16
|
|
model_max_inflight: int = 16
|
|
queue_timeout_seconds: int = 5
|
|
provider_capacity_key: str = ""
|
|
model_capacity_key: str = ""
|
|
attempt_timeout_seconds: int = 600
|
|
runtime_provider: str = ""
|
|
invocation_plan: InvocationPlan | None = None
|
|
supports_tools: bool = False
|
|
|
|
|
|
@dataclass(frozen=True, slots=True)
|
|
class ModelRuntimeSnapshot:
|
|
config_revision: int
|
|
catalog_revision: int
|
|
preparation_id: str
|
|
prepared_snapshot_digest: str
|
|
prepared_input_digest: str
|
|
purpose_attempt_limits: Mapping[str, int]
|
|
purpose_routes: Mapping[str, tuple[ResolvedRoute, ...]]
|
|
purpose_route_call_bounds: Mapping[str, tuple[RouteCallBound, ...]]
|
|
title_start_timeout_seconds: int
|
|
active_run_timeout_seconds: int
|
|
max_run_journal_events: int
|
|
max_run_journal_bytes: int
|
|
|
|
@property
|
|
def main_routes(self) -> tuple[ResolvedRoute, ...]:
|
|
return self.purpose_routes["main_agent"]
|
|
|
|
@property
|
|
def title_route(self) -> ResolvedRoute | None:
|
|
routes = self.purpose_routes.get("title")
|
|
return routes[0] if routes else None
|
|
|
|
|
|
@dataclass(slots=True)
|
|
class _PreparedHandle:
|
|
grant: RoutePreparationGrant
|
|
quote: PreparedRunQuote
|
|
input: AgentInputV3
|
|
host: WebHostContext
|
|
snapshot: ModelRuntimeSnapshot
|
|
tool_registry_payload: Mapping[str, Any]
|
|
secret_fingerprints: Mapping[str, str]
|
|
request_digest: str
|
|
state: str = "PREPARED"
|
|
run: _EvoWebRun | None = None
|
|
start_grant_digest: bytes | None = None
|
|
lock: asyncio.Lock = field(default_factory=asyncio.Lock)
|
|
|
|
|
|
@dataclass(frozen=True, slots=True)
|
|
class _PreparedTombstone:
|
|
request_digest: str
|
|
state: str
|
|
retained_until: int
|
|
|
|
|
|
class RouteHealthBook:
|
|
"""Single-worker circuit breaker with a bounded half-open lease."""
|
|
|
|
def __init__(self) -> None:
|
|
self._revision = -1
|
|
self._policy: Any = None
|
|
self._failures: dict[str, int] = {}
|
|
self._states: dict[str, str] = {}
|
|
self._opened_at: dict[str, float] = {}
|
|
self._half_open_inflight: dict[str, int] = defaultdict(int)
|
|
|
|
def configure(self, config: EvoModelConfig, policy: Any | None = None) -> None:
|
|
if self._revision != config.config_revision:
|
|
self._revision = config.config_revision
|
|
self._policy = policy or config.route_health
|
|
self._failures.clear()
|
|
self._states.clear()
|
|
self._opened_at.clear()
|
|
self._half_open_inflight.clear()
|
|
|
|
def status(self, route_key: str) -> str:
|
|
if self._states.get(route_key) == "open":
|
|
assert self._policy is not None
|
|
if (
|
|
time.monotonic() - self._opened_at.get(route_key, 0)
|
|
>= self._policy.cooldown_seconds
|
|
):
|
|
return "half_open"
|
|
return self._states.get(route_key, "closed")
|
|
|
|
def is_open(self, route_key: str) -> bool:
|
|
return self.status(route_key) == "open"
|
|
|
|
def acquire_half_open(self, route_key: str) -> bool:
|
|
if self.status(route_key) != "half_open":
|
|
return True
|
|
assert self._policy is not None
|
|
if self._half_open_inflight[route_key] >= self._policy.half_open_max_inflight:
|
|
return False
|
|
self._half_open_inflight[route_key] += 1
|
|
return True
|
|
|
|
def release_half_open(self, route_key: str) -> None:
|
|
if self._half_open_inflight.get(route_key, 0) <= 1:
|
|
self._half_open_inflight.pop(route_key, None)
|
|
else:
|
|
self._half_open_inflight[route_key] -= 1
|
|
|
|
def record_success(self, route_key: str) -> None:
|
|
self._failures.pop(route_key, None)
|
|
self._states[route_key] = "closed"
|
|
self._opened_at.pop(route_key, None)
|
|
self._half_open_inflight.pop(route_key, None)
|
|
|
|
def record_failure(self, route_key: str, error_code: str) -> None:
|
|
assert self._policy is not None
|
|
self._half_open_inflight.pop(route_key, None)
|
|
if error_code in self._policy.open_immediately_error_codes:
|
|
self._open(route_key)
|
|
return
|
|
if error_code not in self._policy.counted_error_codes:
|
|
return
|
|
failures = self._failures.get(route_key, 0) + 1
|
|
self._failures[route_key] = failures
|
|
if failures >= self._policy.failure_threshold:
|
|
self._open(route_key)
|
|
|
|
def _open(self, route_key: str) -> None:
|
|
self._states[route_key] = "open"
|
|
self._opened_at[route_key] = time.monotonic()
|
|
|
|
|
|
class _SmoothWeightedRoundRobin:
|
|
def __init__(self) -> None:
|
|
self._revision = -1
|
|
self._current: dict[tuple[str, str], int] = {}
|
|
self._locks: dict[str, asyncio.Lock] = {}
|
|
|
|
def configure(self, revision: int) -> None:
|
|
if revision != self._revision:
|
|
self._revision = revision
|
|
self._current.clear()
|
|
self._locks.clear()
|
|
|
|
async def choose(
|
|
self, pool_id: str, weighted_names: Sequence[tuple[str, int]]
|
|
) -> str:
|
|
lock = self._locks.setdefault(pool_id, asyncio.Lock())
|
|
async with lock:
|
|
total = sum(weight for _, weight in weighted_names)
|
|
best_name = ""
|
|
best_value: int | None = None
|
|
for name, weight in weighted_names:
|
|
key = (pool_id, name)
|
|
current = self._current.get(key, 0) + weight
|
|
self._current[key] = current
|
|
if best_value is None or current > best_value:
|
|
best_name, best_value = name, current
|
|
self._current[(pool_id, best_name)] -= total
|
|
return best_name
|
|
|
|
|
|
class InvocationCapacityBook:
|
|
"""Process-local Provider and Model concurrency gates for Web runtime calls."""
|
|
|
|
def __init__(self) -> None:
|
|
self._semaphores: dict[tuple[str, int], asyncio.BoundedSemaphore] = {}
|
|
|
|
def _semaphore(self, key: str, limit: int) -> asyncio.BoundedSemaphore:
|
|
return self._semaphores.setdefault(
|
|
(key, limit), asyncio.BoundedSemaphore(limit)
|
|
)
|
|
|
|
async def acquire(self, model: Any) -> tuple[asyncio.BoundedSemaphore, ...]:
|
|
metadata = getattr(model, "metadata", None) or {}
|
|
provider_key = str(metadata.get("capacity_provider_key") or "")
|
|
model_key = str(metadata.get("capacity_model_key") or "")
|
|
if not provider_key or not model_key:
|
|
return ()
|
|
provider = self._semaphore(
|
|
provider_key, int(metadata.get("provider_max_inflight") or 16)
|
|
)
|
|
model_gate = self._semaphore(
|
|
model_key, int(metadata.get("model_max_inflight") or 16)
|
|
)
|
|
acquired: list[asyncio.BoundedSemaphore] = []
|
|
try:
|
|
async with asyncio.timeout(
|
|
int(metadata.get("capacity_queue_timeout_seconds") or 5)
|
|
):
|
|
await provider.acquire()
|
|
acquired.append(provider)
|
|
await model_gate.acquire()
|
|
acquired.append(model_gate)
|
|
except TimeoutError as exc:
|
|
for gate in reversed(acquired):
|
|
gate.release()
|
|
raise EvoRuntimeError("MODEL_CAPACITY_EXHAUSTED") from exc
|
|
return tuple(acquired)
|
|
|
|
@staticmethod
|
|
def release(lease: tuple[asyncio.BoundedSemaphore, ...]) -> None:
|
|
for gate in reversed(lease):
|
|
gate.release()
|
|
|
|
|
|
class EvoModelRuntime:
|
|
"""Prepare, authorize and execute one frozen Evo Agent transaction."""
|
|
|
|
def __init__(
|
|
self,
|
|
store: FileEvoModelConfigStore,
|
|
*,
|
|
admission_verifier: AdmissionGrantVerifier,
|
|
quote_authority: HmacGrantAuthority | None = None,
|
|
identity_key_ring: HmacKeyRing | None = None,
|
|
secret_resolver: SecretResolver | None = None,
|
|
model_factory: Callable[..., Any] | None = None,
|
|
agent_factory: Callable[
|
|
[ModelRuntimeSnapshot, WebHostContext, AgentModelSet], Any
|
|
]
|
|
| None = None,
|
|
runtime_instance_id: str | None = None,
|
|
host_registry: HostExecutionRegistry | None = None,
|
|
) -> None:
|
|
self.store = store
|
|
self.admission_verifier = admission_verifier
|
|
if quote_authority is None and isinstance(
|
|
admission_verifier, HmacGrantAuthority
|
|
):
|
|
quote_authority = admission_verifier
|
|
if quote_authority is None:
|
|
raise ValueError("quote authority is required")
|
|
if identity_key_ring is None:
|
|
raise ValueError("config identity key ring is required")
|
|
self.quote_authority = quote_authority
|
|
self.identity_key_ring = identity_key_ring
|
|
self.secret_resolver = secret_resolver
|
|
self.model_factory = model_factory or self._default_model_factory
|
|
self.agent_factory = agent_factory or self._default_agent_factory
|
|
self.runtime_instance_id = runtime_instance_id or str(uuid.uuid4())
|
|
self.host_registry = host_registry
|
|
self._route_health = RouteHealthBook()
|
|
self._provider_health = RouteHealthBook()
|
|
self._pool = _SmoothWeightedRoundRobin()
|
|
self._capacity = InvocationCapacityBook()
|
|
self._prepared: dict[str, _PreparedHandle] = {}
|
|
self._prepared_requests: dict[tuple[str, str, str], str] = {}
|
|
self._prepared_tombstones: dict[tuple[str, str, str], _PreparedTombstone] = {}
|
|
self._started_grants: dict[str, tuple[bytes, _EvoWebRun]] = {}
|
|
self._registry_lock = asyncio.Lock()
|
|
|
|
def _record_runtime_observation(
|
|
self,
|
|
route: ResolvedRoute,
|
|
*,
|
|
purpose: str,
|
|
outcome: str,
|
|
error_code: str | None,
|
|
) -> None:
|
|
"""Record call metadata when the configured store supports telemetry.
|
|
|
|
Observability must never change the outcome of a request, including
|
|
when the telemetry database is unavailable.
|
|
"""
|
|
|
|
recorder = getattr(self.store, "record_runtime_observation", None)
|
|
if not callable(recorder):
|
|
return
|
|
params = (
|
|
route.invocation_plan.sdk_params
|
|
if route.invocation_plan is not None
|
|
else {}
|
|
)
|
|
extra_body = params.get("extra_body")
|
|
if not isinstance(extra_body, Mapping):
|
|
extra_body = {}
|
|
response_format = params.get("response_format")
|
|
strategy = {
|
|
"enable_thinking": extra_body.get("enable_thinking"),
|
|
"structured_output": bool(response_format),
|
|
"tool_choice": params.get("tool_choice"),
|
|
}
|
|
try:
|
|
recorder(
|
|
config_revision=route.identity.config_revision,
|
|
provider_id=route.identity.provider_id,
|
|
model_profile_id=route.ref.model,
|
|
provider_model_id=route.identity.model_id,
|
|
api_mode=route.identity.api_mode,
|
|
purpose=purpose,
|
|
outcome=outcome,
|
|
error_code=error_code,
|
|
strategy=strategy,
|
|
)
|
|
except Exception:
|
|
# Telemetry is intentionally non-blocking and non-authoritative.
|
|
return
|
|
|
|
async def prepare_model_run(
|
|
self,
|
|
grant: RoutePreparationGrant,
|
|
agent_input: AgentInputV3,
|
|
host: WebHostContext,
|
|
) -> PreparedRunQuote:
|
|
self.admission_verifier.require_preparation(grant)
|
|
self._validate_preparation_echo(grant, agent_input)
|
|
await self._validate_continuation(grant, agent_input, host)
|
|
tool_payload = self._tool_registry_payload(host)
|
|
request_key = (grant.subject_id, grant.request_id, grant.turn_id)
|
|
request_digest = self._preparation_request_digest(
|
|
grant, agent_input, tool_payload
|
|
)
|
|
async with self._registry_lock:
|
|
self._expire_prepared_locked()
|
|
replay = self._prepared_replay_locked(request_key, request_digest)
|
|
if replay is not None:
|
|
return replay
|
|
config = self.store.load()
|
|
self._route_health.configure(config)
|
|
self._provider_health.configure(
|
|
config, config.provider_health or config.route_health
|
|
)
|
|
self._pool.configure(config.config_revision)
|
|
if not self.identity_key_ring.contains(config.config_identity_key_id):
|
|
raise EvoRuntimeError("CONFIG_IDENTITY_KEY_UNKNOWN")
|
|
main_selector = config.resolve_main_selector(grant.requested_model_ref)
|
|
main_model = config.providers[main_selector.provider].models[
|
|
main_selector.model
|
|
]
|
|
main_provider = config.providers[main_selector.provider]
|
|
current_options_schema_hash = model_options_schema_hash(
|
|
model_profile_id=grant.requested_model_ref,
|
|
user_options=main_model.user_options,
|
|
supports_reasoning=main_model.supports_reasoning,
|
|
reasoning_mode=main_model.reasoning_mode,
|
|
allowed_reasoning_efforts=main_model.allowed_reasoning_efforts,
|
|
parameter_constraints=main_model.parameter_constraints,
|
|
adapter_id=main_provider.adapter_id,
|
|
adapter_revision=main_provider.adapter_revision,
|
|
)
|
|
supplied_options_schema_hash = str(
|
|
agent_input.metadata.get("model_options_schema_hash") or ""
|
|
)
|
|
if (
|
|
supplied_options_schema_hash
|
|
and supplied_options_schema_hash != current_options_schema_hash
|
|
):
|
|
raise EvoRuntimeError("MODEL_OPTIONS_STALE")
|
|
supplied_model_options = dict(agent_input.metadata.get("model_options") or {})
|
|
combined_user_options = dict(supplied_model_options)
|
|
if main_model.supports_reasoning:
|
|
combined_user_options["reasoning"] = (
|
|
"off"
|
|
if grant.reasoning_effort == "disabled"
|
|
else grant.reasoning_effort
|
|
)
|
|
validated_user_options = validate_user_model_options(
|
|
supplied=combined_user_options,
|
|
user_options=main_model.user_options,
|
|
supports_reasoning=main_model.supports_reasoning,
|
|
reasoning_mode=main_model.reasoning_mode,
|
|
allowed_reasoning_efforts=main_model.allowed_reasoning_efforts,
|
|
default_reasoning_effort=str(
|
|
main_model.reasoning_enabled_params.get("reasoning") or ""
|
|
),
|
|
parameter_constraints=main_model.parameter_constraints,
|
|
)
|
|
validated_user_options.pop("reasoning", None)
|
|
async with self._registry_lock:
|
|
per_subject = sum(
|
|
handle.state == "PREPARED"
|
|
and handle.grant.subject_id == grant.subject_id
|
|
for handle in self._prepared.values()
|
|
)
|
|
active_total = sum(
|
|
handle.state == "PREPARED" for handle in self._prepared.values()
|
|
)
|
|
if (
|
|
per_subject >= config.web_runtime.max_prepared_runs_per_subject
|
|
or active_total >= config.web_runtime.max_prepared_runs_total
|
|
):
|
|
raise EvoRuntimeError("PREPARATION_CAPACITY_EXCEEDED")
|
|
media_token_bound = 0
|
|
for item in agent_input.media:
|
|
if hasattr(item, "token_bound"):
|
|
bound = int(item.token_bound)
|
|
elif isinstance(item, Mapping) and "token_bound" in item:
|
|
bound = int(item["token_bound"])
|
|
else:
|
|
raise EvoRuntimeError("TOKEN_BOUND_UNAVAILABLE")
|
|
if bound < 0:
|
|
raise EvoRuntimeError("TOKEN_BOUND_UNAVAILABLE")
|
|
media_token_bound += bound
|
|
tool_snapshot_id = self.quote_authority.tool_registry_snapshot_id(tool_payload)
|
|
prepared_input_payload = {
|
|
"agent_input": agent_input.projection(),
|
|
"checkpoint_snapshot_id": grant.checkpoint_snapshot_id,
|
|
"tool_registry_snapshot_id": tool_snapshot_id,
|
|
}
|
|
prepared_input_digest = self.quote_authority.prepared_input_digest(
|
|
prepared_input_payload
|
|
)
|
|
enabled_purposes = _enabled_purposes(grant.title_policy)
|
|
purpose_routes = await self._freeze_purpose_routes(
|
|
config,
|
|
grant,
|
|
validated_user_options,
|
|
)
|
|
purpose_attempt_limits = {
|
|
purpose: config.purpose_call_limits[purpose].max_attempts_per_run
|
|
for purpose in enabled_purposes
|
|
}
|
|
bounds = self._build_call_bounds(purpose_routes)
|
|
purpose_routes = self._compile_purpose_routes(
|
|
purpose_routes,
|
|
bounds,
|
|
grant.reasoning_effort,
|
|
)
|
|
initial_input_bound = (
|
|
_payload_token_bound(
|
|
{"input": agent_input.projection(), "tool_registry": tool_payload}
|
|
)
|
|
+ media_token_bound
|
|
)
|
|
if any(
|
|
initial_input_bound > bound.payload_input_hard_cap
|
|
for purpose_bounds in bounds.values()
|
|
for bound in purpose_bounds
|
|
):
|
|
raise EvoRuntimeError("MODEL_CONTEXT_WINDOW_EXCEEDED")
|
|
reserve = self._run_reserve(purpose_attempt_limits, bounds)
|
|
total_attempts = sum(purpose_attempt_limits.values())
|
|
public_routes = {
|
|
purpose: {
|
|
"primary": routes[0].identity,
|
|
"fallbacks": tuple(route.identity for route in routes[1:]),
|
|
}
|
|
for purpose, routes in purpose_routes.items()
|
|
}
|
|
public_invocation_plans = {
|
|
purpose: {
|
|
"primary": routes[0].invocation_plan.projection()
|
|
if routes[0].invocation_plan
|
|
else {},
|
|
"fallbacks": tuple(
|
|
route.invocation_plan.projection() if route.invocation_plan else {}
|
|
for route in routes[1:]
|
|
),
|
|
}
|
|
for purpose, routes in purpose_routes.items()
|
|
}
|
|
public_bounds = {
|
|
purpose: tuple(bounds[purpose]) for purpose in purpose_attempt_limits
|
|
}
|
|
quotes = {
|
|
route.quote.quote_id: route.quote
|
|
for routes in purpose_routes.values()
|
|
for route in routes
|
|
}
|
|
execution_id = str(uuid.uuid4())
|
|
snapshot_payload = {
|
|
"execution_id": execution_id,
|
|
"predecessor_execution_id": grant.predecessor_execution_id,
|
|
"predecessor_checkpoint_id": grant.predecessor_checkpoint_id,
|
|
"predecessor_owner_epoch": grant.predecessor_owner_epoch,
|
|
"continuation_pending_hash": grant.continuation_pending_hash,
|
|
"continuation_decision_hash": grant.continuation_decision_hash,
|
|
"request_id": grant.request_id,
|
|
"turn_id": grant.turn_id,
|
|
"thread_id": grant.thread_id,
|
|
"subject_id": grant.subject_id,
|
|
"requested_model_ref": grant.requested_model_ref,
|
|
"config_revision": config.config_revision,
|
|
"gateway_input_digest": grant.gateway_input_digest,
|
|
"prepared_input_digest": prepared_input_digest,
|
|
"checkpoint_snapshot_id": grant.checkpoint_snapshot_id,
|
|
"tool_registry_snapshot_id": tool_snapshot_id,
|
|
"purpose_routes": public_routes,
|
|
"purpose_invocation_plans": public_invocation_plans,
|
|
"purpose_route_call_bounds": public_bounds,
|
|
"purpose_attempt_limits": purpose_attempt_limits,
|
|
"provider_run_reserve_microunits": reserve,
|
|
"route_health": asdict(config.route_health),
|
|
}
|
|
prepared_snapshot_digest = self.quote_authority.prepared_snapshot_digest(
|
|
snapshot_payload
|
|
)
|
|
preparation_id = str(uuid.uuid4())
|
|
issued_at = now_ms()
|
|
quote = self.quote_authority.sign_quote(
|
|
execution_id=execution_id,
|
|
preparation_id=preparation_id,
|
|
request_id=grant.request_id,
|
|
turn_id=grant.turn_id,
|
|
thread_id=grant.thread_id,
|
|
subject_id=grant.subject_id,
|
|
requested_model_ref=grant.requested_model_ref,
|
|
plan=grant.plan,
|
|
roles=grant.roles,
|
|
requires_vision=grant.requires_vision,
|
|
reasoning_effort=grant.reasoning_effort,
|
|
title_policy=grant.title_policy,
|
|
gateway_input_digest=grant.gateway_input_digest,
|
|
prepared_snapshot_digest=prepared_snapshot_digest,
|
|
prepared_input_digest=prepared_input_digest,
|
|
config_revision=config.config_revision,
|
|
catalog_revision=config.catalog_revision,
|
|
enabled_purposes=enabled_purposes,
|
|
purpose_routes=public_routes,
|
|
purpose_route_call_bounds=bounds,
|
|
purpose_attempt_limits=purpose_attempt_limits,
|
|
total_max_attempts=total_attempts,
|
|
pricing_quotes=quotes,
|
|
quote_ids=tuple(sorted(quotes)),
|
|
provider_run_reserve_microunits=reserve,
|
|
checkpoint_snapshot_id=grant.checkpoint_snapshot_id,
|
|
tool_registry_snapshot_id=tool_snapshot_id,
|
|
turn_fencing_token=grant.turn_fencing_token,
|
|
route_semantics_hashes=tuple(
|
|
sorted(
|
|
{
|
|
route.identity.route_semantics_hash
|
|
for routes in purpose_routes.values()
|
|
for route in routes
|
|
}
|
|
)
|
|
),
|
|
issued_at=issued_at,
|
|
expires_at=min(
|
|
grant.expires_at,
|
|
issued_at + config.web_runtime.prepare_ttl_seconds * 1000,
|
|
),
|
|
)
|
|
snapshot = ModelRuntimeSnapshot(
|
|
config_revision=config.config_revision,
|
|
catalog_revision=config.catalog_revision,
|
|
preparation_id=preparation_id,
|
|
prepared_snapshot_digest=prepared_snapshot_digest,
|
|
prepared_input_digest=prepared_input_digest,
|
|
purpose_attempt_limits=purpose_attempt_limits,
|
|
purpose_routes=purpose_routes,
|
|
purpose_route_call_bounds=public_bounds,
|
|
title_start_timeout_seconds=config.web_runtime.title_start_timeout_seconds,
|
|
active_run_timeout_seconds=config.web_runtime.active_run_timeout_seconds,
|
|
max_run_journal_events=config.web_runtime.max_run_journal_events,
|
|
max_run_journal_bytes=config.web_runtime.max_run_journal_bytes,
|
|
)
|
|
fingerprints = {
|
|
secret_ref: fingerprint
|
|
for routes in purpose_routes.values()
|
|
for route in routes
|
|
for secret_ref, fingerprint in route.secret_fingerprints.items()
|
|
}
|
|
handle = _PreparedHandle(
|
|
grant=grant,
|
|
quote=quote,
|
|
input=agent_input,
|
|
host=host,
|
|
snapshot=snapshot,
|
|
tool_registry_payload=tool_payload,
|
|
secret_fingerprints=fingerprints,
|
|
request_digest=request_digest,
|
|
)
|
|
async with self._registry_lock:
|
|
self._expire_prepared_locked()
|
|
replay = self._prepared_replay_locked(request_key, request_digest)
|
|
if replay is not None:
|
|
return replay
|
|
if grant.grant_id in (
|
|
item.grant.grant_id for item in self._prepared.values()
|
|
):
|
|
raise EvoRuntimeError("PREPARATION_CONFLICT")
|
|
per_subject = sum(
|
|
item.state == "PREPARED" and item.grant.subject_id == grant.subject_id
|
|
for item in self._prepared.values()
|
|
)
|
|
active_total = sum(
|
|
item.state == "PREPARED" for item in self._prepared.values()
|
|
)
|
|
if (
|
|
per_subject >= config.web_runtime.max_prepared_runs_per_subject
|
|
or active_total >= config.web_runtime.max_prepared_runs_total
|
|
):
|
|
raise EvoRuntimeError("PREPARATION_CAPACITY_EXCEEDED")
|
|
self._prepared[preparation_id] = handle
|
|
self._prepared_requests[request_key] = preparation_id
|
|
return quote
|
|
|
|
def inspect(self, execution_id: str, *, owner_epoch: int | None = None,
|
|
boot_id: str | None = None) -> dict:
|
|
if self.host_registry is None:
|
|
return {"execution_id": execution_id, "status": "unknown",
|
|
"resources_confirmed_exited": False, "source": "registry_disabled"}
|
|
state = self.host_registry.inspect_control(
|
|
execution_id, owner_epoch=owner_epoch, boot_id=boot_id)
|
|
run = next((h.run for h in self._prepared.values()
|
|
if h.run is not None and h.run.run_id == execution_id), None)
|
|
if run is not None:
|
|
state["terminal_error_code"] = run._terminal_error_code
|
|
return state
|
|
|
|
async def cancel(self, execution_id: str, *, reason: str,
|
|
owner_epoch: int | None = None, boot_id: str | None = None) -> str:
|
|
run = next((run for _, run in self._started_grants.values()
|
|
if run.run_id == execution_id), None)
|
|
if run is None:
|
|
raise EvoRuntimeError("EXECUTION_UNKNOWN")
|
|
if self.host_registry is not None:
|
|
self.host_registry.accept_cancel(
|
|
execution_id, owner_epoch=owner_epoch, boot_id=boot_id, reason=reason)
|
|
# Consume the committed intent before yielding. Later transfers do not
|
|
# revoke already accepted cancellation; the run owns its cleanup task.
|
|
run._request_cancel(reason)
|
|
return await run.wait_stopped(timeout=2.0)
|
|
|
|
async def recover_terminal(self, execution_id: str, *, sink: Any) -> str:
|
|
"""Explicit host repair of durable cleanup evidence; never executes Graph."""
|
|
registry = self.host_registry
|
|
intent = registry.terminal_intent(execution_id) if registry is not None else None
|
|
if registry is None or intent is None:
|
|
raise EvoRuntimeError("TERMINAL_CLEANUP_EVIDENCE_REQUIRED")
|
|
event = EvoRuntimeEvent(**intent['event'])
|
|
digest = hashlib.sha256(canonical_json_v1(asdict(event))).hexdigest()
|
|
if digest != intent['digest']:
|
|
raise EvoRuntimeError("EVO_EVENT_CONFLICT")
|
|
if intent['phase'] == 'prepared':
|
|
try:
|
|
confirmation = await sink.confirm(event.event_id, digest)
|
|
if confirmation == 'absent':
|
|
await sink.commit(event)
|
|
confirmation = await sink.confirm(event.event_id, digest)
|
|
except Exception as exc:
|
|
raise EvoRuntimeError("EVENT_COMMIT_INDETERMINATE") from exc
|
|
if confirmation != 'committed':
|
|
raise EvoRuntimeError("EVO_EVENT_CONFLICT" if confirmation == 'conflict'
|
|
else "EVENT_COMMIT_INDETERMINATE")
|
|
await asyncio.to_thread(registry.confirm_terminal, execution_id, digest=digest)
|
|
outcome = str(event.payload['outcome'])
|
|
await asyncio.to_thread(registry.finish, execution_id, outcome=outcome,
|
|
checkpoint_id=str(event.payload.get('checkpoint_id') or ''))
|
|
run = next((h.run for h in self._prepared.values()
|
|
if h.run is not None and h.run.run_id == execution_id), None)
|
|
if run is not None and run._terminal_task is not None and run._terminal_task.done():
|
|
async with run._condition:
|
|
if not any(item.event_id == event.event_id for item in run._journal):
|
|
run._journal.append(event)
|
|
run._journal_bytes += len(canonical_json_v1(asdict(event)))
|
|
run._sequence = max(run._sequence, event.sequence)
|
|
run._terminal_event = event
|
|
run._state = 'TERMINAL'
|
|
run._terminal_error_code = None
|
|
run._condition.notify_all()
|
|
if run._prepared.state == 'CONSTRUCTION_FAILED':
|
|
async with self._registry_lock:
|
|
self._retire_prepared_locked(run._prepared)
|
|
return outcome
|
|
|
|
async def _checkpoint_store_id(self, host):
|
|
# No portable durable identity exists in BaseCheckpointSaver. Only the
|
|
# inspected SQLite implementation is supported by this opt-in protocol.
|
|
from langgraph.checkpoint.sqlite.aio import AsyncSqliteSaver
|
|
import os
|
|
if type(host.checkpointer) is not AsyncSqliteSaver:
|
|
raise EvoRuntimeError("CHECKPOINT_STORE_UNSUPPORTED")
|
|
async with host.checkpointer.conn.execute("PRAGMA database_list") as cursor:
|
|
databases = await cursor.fetchall()
|
|
filename = next((row[2] for row in databases if row[1] == "main"), "")
|
|
if not filename:
|
|
raise EvoRuntimeError("CHECKPOINT_STORE_UNSUPPORTED")
|
|
info = os.stat(filename)
|
|
return f"sqlite:{info.st_dev}:{info.st_ino}"
|
|
|
|
async def _validate_continuation(self, grant, agent_input, host):
|
|
if self.host_registry is None:
|
|
return
|
|
from langgraph.types import Command
|
|
command = agent_input.message
|
|
if not grant.predecessor_execution_id:
|
|
if isinstance(command, Command):
|
|
raise EvoRuntimeError("CONTINUATION_REQUIRED")
|
|
return
|
|
pending = self.host_registry.continuation(grant.predecessor_execution_id)
|
|
if (pending["owner_epoch"] != grant.predecessor_owner_epoch
|
|
or isinstance(grant.predecessor_owner_epoch, bool)
|
|
or pending["pending_hash"] != grant.continuation_pending_hash
|
|
or pending["checkpoint_id"] != grant.predecessor_checkpoint_id
|
|
or grant.checkpoint_snapshot_id != grant.predecessor_checkpoint_id):
|
|
raise EvoRuntimeError("CONTINUATION_AUTHORIZATION_INVALID")
|
|
if not isinstance(command, Command) or command.update or command.goto or command.graph:
|
|
raise EvoRuntimeError("CONTINUATION_DECISION_INVALID")
|
|
resume = command.resume
|
|
if (not isinstance(resume, dict) or set(resume) != {"decisions"}
|
|
or hashlib.sha256(canonical_json_v1(resume)).hexdigest() != grant.continuation_decision_hash):
|
|
raise EvoRuntimeError("CONTINUATION_DECISION_INVALID")
|
|
intent = self.host_registry.terminal_intent(grant.predecessor_execution_id)
|
|
if intent is None or intent["phase"] != "registry_finished":
|
|
raise EvoRuntimeError("CONTINUATION_INVALID")
|
|
payload = intent["event"]["payload"]
|
|
identity = {k: payload[k] for k in (
|
|
"checkpoint_thread_id", "checkpoint_id", "checkpoint_ns", "pending_interrupts")}
|
|
if identity["checkpoint_thread_id"] != agent_input.checkpoint_thread_id or identity["checkpoint_ns"] != "":
|
|
raise EvoRuntimeError("CONTINUATION_INVALID")
|
|
checkpoint = await host.checkpointer.aget_tuple({"configurable": {
|
|
"thread_id": agent_input.checkpoint_thread_id, "checkpoint_ns": identity["checkpoint_ns"]}})
|
|
if (checkpoint is None or checkpoint.config["configurable"]["checkpoint_id"] != identity["checkpoint_id"]
|
|
or checkpoint.checkpoint["channel_values"].get("_verified_review_mode", {}).get("mode") != "manual"):
|
|
raise EvoRuntimeError("CONTINUATION_CHECKPOINT_STALE")
|
|
actual = []
|
|
for _, channel, value in checkpoint.pending_writes or ():
|
|
if channel == "__interrupt__":
|
|
for item in value:
|
|
projected = {"id": item.id, "value": item.value}
|
|
if projected not in actual:
|
|
actual.append(projected)
|
|
if actual != identity["pending_interrupts"] or not actual:
|
|
raise EvoRuntimeError("CONTINUATION_PENDING_STALE")
|
|
actions = [action for item in actual for action in item["value"].get("action_requests", [])]
|
|
decisions = resume["decisions"]
|
|
if (len(actual) != 1 or not actions or not isinstance(decisions, list)
|
|
or len(decisions) != len(actions)
|
|
or any(not isinstance(d, dict) or d.get("type") not in {"approve", "reject"}
|
|
or set(d) - {"type", "message"} for d in decisions)):
|
|
raise EvoRuntimeError("CONTINUATION_DECISION_INVALID")
|
|
|
|
async def start_web_run(self, admission: AdmissionGrant) -> EvoWebRun:
|
|
self.admission_verifier.require_admission(admission)
|
|
digest = canonical_json_v1(admission.unsigned_payload())
|
|
async with self._registry_lock:
|
|
replay = self._started_grants.get(admission.grant_id)
|
|
if replay is not None:
|
|
if replay[0] != digest:
|
|
raise EvoRuntimeError("CONTRACT_REPLAYED")
|
|
return replay[1]
|
|
if self.host_registry is not None:
|
|
binding = self.host_registry.lookup_grant(
|
|
admission.grant_id, hashlib.sha256(digest).hexdigest()
|
|
)
|
|
if binding is not None:
|
|
if self.host_registry.inspect(binding["execution_id"])["status"] == "failed":
|
|
raise EvoRuntimeError("EXECUTION_FAILED_RETRY_REQUIRES_NEW_ADMISSION")
|
|
raise EvoRuntimeError("EXECUTION_UNKNOWN")
|
|
handle = self._prepared.get(admission.preparation_id)
|
|
if handle is None:
|
|
raise EvoRuntimeError("RUN_LOST")
|
|
async with handle.lock:
|
|
if handle.state == "STARTED" and handle.run is not None:
|
|
return handle.run
|
|
if handle.state != "PREPARED":
|
|
raise EvoRuntimeError("CONTRACT_REPLAYED")
|
|
if handle.quote.expires_at < now_ms():
|
|
handle.state = "EXPIRED"
|
|
raise EvoRuntimeError("PREPARATION_STALE")
|
|
self._validate_admission_echo(handle.quote, admission)
|
|
self._validate_still_fresh(handle)
|
|
if handle.host.runtime_event_sink is None:
|
|
raise EvoRuntimeError("EVENT_INGRESS_UNAVAILABLE")
|
|
run = _EvoWebRun(
|
|
runtime=self,
|
|
admission=admission,
|
|
prepared=handle,
|
|
agent=None,
|
|
model_set=AgentModelSet(None, None, None),
|
|
)
|
|
if self.host_registry is not None:
|
|
store_id = await self._checkpoint_store_id(handle.host)
|
|
self.host_registry.claim_checkpoint(
|
|
run.run_id, store_id=store_id,
|
|
checkpoint_thread_id=handle.input.checkpoint_thread_id,
|
|
subject_id=admission.subject_id)
|
|
try:
|
|
await self._validate_continuation(handle.grant, handle.input, handle.host)
|
|
if self.host_registry is not None:
|
|
self.host_registry.bind(
|
|
execution_id=run.run_id, grant_id=admission.grant_id,
|
|
digest=hashlib.sha256(digest).hexdigest(),
|
|
thread_id=admission.thread_id, turn_id=admission.turn_id,
|
|
predecessor_execution_id=handle.grant.predecessor_execution_id,
|
|
predecessor_checkpoint_id=handle.grant.predecessor_checkpoint_id,
|
|
predecessor_owner_epoch=handle.grant.predecessor_owner_epoch,
|
|
continuation_pending_hash=handle.grant.continuation_pending_hash,
|
|
continuation_decision_hash=handle.grant.continuation_decision_hash,
|
|
)
|
|
except BaseException:
|
|
if self.host_registry is not None:
|
|
self.host_registry.release_unbound_checkpoint(run.run_id)
|
|
raise
|
|
token = _construction_owner.set(run)
|
|
try:
|
|
model_set = self._build_model_set(
|
|
handle.snapshot, admission.reasoning_effort
|
|
)
|
|
run._model_set = model_set
|
|
run._agent = self.agent_factory(handle.snapshot, handle.host, model_set)
|
|
except BaseException:
|
|
# Retain the owner on the handle even if rollback itself fails.
|
|
handle.run = run
|
|
handle.state = "CONSTRUCTION_FAILED"
|
|
run._start_terminal("failed", error_code="RUN_CONSTRUCTION_FAILED")
|
|
await run._ensure_cleanup()
|
|
await run.wait_stopped(timeout=0.25)
|
|
raise
|
|
finally:
|
|
_construction_owner.reset(token)
|
|
handle.state = "STARTED"
|
|
handle.run = run
|
|
handle.start_grant_digest = digest
|
|
async with self._registry_lock:
|
|
self._started_grants[admission.grant_id] = (digest, run)
|
|
run._start()
|
|
return run
|
|
|
|
async def cancel_prepared_run(self, preparation_id: str, *, reason: str) -> bool:
|
|
"""Idempotently invalidate a prepared handle that never reached Start."""
|
|
|
|
_ = reason
|
|
async with self._registry_lock:
|
|
handle = self._prepared.get(str(preparation_id))
|
|
if handle is None:
|
|
return False
|
|
async with handle.lock:
|
|
if handle.state in {"CANCELLED", "EXPIRED"}:
|
|
return True
|
|
if handle.state != "PREPARED":
|
|
return False
|
|
handle.state = "CANCELLED"
|
|
async with self._registry_lock:
|
|
self._retire_prepared_locked(handle)
|
|
return True
|
|
|
|
async def get_catalog(self, subject: VerifiedModelSubject) -> ModelCatalog:
|
|
if not self.admission_verifier.verify_subject(subject):
|
|
raise EvoRuntimeError("CATALOG_SUBJECT_INVALID")
|
|
config = self.store.load()
|
|
entries: list[ModelCatalogEntry] = []
|
|
for alias, selector_id in config.main_routes.selectable.items():
|
|
selector = config.route_selectors[selector_id]
|
|
model = config.providers[selector.provider].models[selector.model]
|
|
alias_config = config.aliases.get(alias)
|
|
if self._access_allowed(
|
|
model, subject.plan, subject.roles, subject.subject_id
|
|
) and (
|
|
alias_config is None
|
|
or self._access_policy_allowed(
|
|
alias_config.access, subject.roles, subject.subject_id
|
|
)
|
|
):
|
|
health_states = [
|
|
self._route_health.status(route.key())
|
|
for route in config.concrete_routes(selector_id)
|
|
]
|
|
aggregate = (
|
|
"closed"
|
|
if "closed" in health_states
|
|
else ("half_open" if "half_open" in health_states else "open")
|
|
)
|
|
entries.append(
|
|
# Alias defaults are display defaults only. The runtime
|
|
# already compiles them into the frozen route.
|
|
ModelCatalogEntry(
|
|
alias=alias,
|
|
provider=selector.provider,
|
|
supports_vision=model.supports_vision,
|
|
supports_reasoning=model.supports_reasoning,
|
|
allowed_reasoning_efforts=model.allowed_reasoning_efforts,
|
|
context_window=model.context_window,
|
|
max_output_tokens=model.max_output_tokens,
|
|
reasoning_mode=model.reasoning_mode,
|
|
default_reasoning_effort=str(
|
|
model.reasoning_enabled_params.get("reasoning") or ""
|
|
),
|
|
billing_sku=model.billing_sku,
|
|
quote=model.quote,
|
|
health=aggregate,
|
|
display_name=(
|
|
alias_config.display_name if alias_config else alias
|
|
),
|
|
provider_display_name=config.providers[
|
|
selector.provider
|
|
].display_name,
|
|
description=model.description,
|
|
version_policy=model.version_policy,
|
|
resolved_model_revision=model.resolved_model_revision,
|
|
reproducible=model.reproducible,
|
|
capabilities=tuple(
|
|
key
|
|
for key, enabled in model.capabilities.items()
|
|
if enabled
|
|
),
|
|
user_options={
|
|
name: {
|
|
**dict(option),
|
|
**(
|
|
{"default": alias_config.defaults[name]}
|
|
if alias_config is not None
|
|
and name in alias_config.defaults
|
|
else {"default": model.params[name]}
|
|
if name in model.params
|
|
else {}
|
|
),
|
|
}
|
|
for name, option in model.user_options.items()
|
|
},
|
|
parameter_constraints=model.parameter_constraints,
|
|
options_schema_hash=model_options_schema_hash(
|
|
model_profile_id=alias,
|
|
user_options=model.user_options,
|
|
supports_reasoning=model.supports_reasoning,
|
|
reasoning_mode=model.reasoning_mode,
|
|
allowed_reasoning_efforts=model.allowed_reasoning_efforts,
|
|
parameter_constraints=model.parameter_constraints,
|
|
adapter_id=config.providers[selector.provider].adapter_id,
|
|
adapter_revision=config.providers[
|
|
selector.provider
|
|
].adapter_revision,
|
|
),
|
|
)
|
|
)
|
|
allowed_aliases = {entry.alias for entry in entries}
|
|
default_alias = config.main_routes.default_alias
|
|
if default_alias not in allowed_aliases:
|
|
default_alias = min(allowed_aliases) if allowed_aliases else ""
|
|
return ModelCatalog(config.catalog_revision, default_alias, tuple(entries))
|
|
|
|
def turn_lease_ttl_seconds(self) -> int:
|
|
config = self.store.load()
|
|
return (
|
|
config.web_runtime.prepare_ttl_seconds
|
|
+ config.web_runtime.turn_lease_grace_seconds
|
|
)
|
|
|
|
async def _freeze_purpose_routes(
|
|
self,
|
|
config: EvoModelConfig,
|
|
grant: RoutePreparationGrant,
|
|
model_options: Mapping[str, Any] | None = None,
|
|
) -> Mapping[str, tuple[ResolvedRoute, ...]]:
|
|
main_selector = config.resolve_main_selector(grant.requested_model_ref)
|
|
main_refs = [await self._select_concrete(config, main_selector.selector_id)]
|
|
for fallback_selector in config.fallback_selectors(main_selector.selector_id):
|
|
main_refs.append(await self._select_concrete(config, fallback_selector))
|
|
main_model = config.route_model(main_refs[0])
|
|
alias_config = config.aliases.get(main_selector.alias)
|
|
if not self._access_allowed(
|
|
main_model, grant.plan, grant.roles, grant.subject_id
|
|
) or (
|
|
alias_config is not None
|
|
and not self._access_policy_allowed(
|
|
alias_config.access, grant.roles, grant.subject_id
|
|
)
|
|
):
|
|
raise EvoRuntimeError("MODEL_ACCESS_DENIED")
|
|
if grant.requires_vision and not main_model.supports_vision:
|
|
raise EvoRuntimeError("MODEL_CAPABILITY_UNAVAILABLE")
|
|
if grant.reasoning_effort != "disabled":
|
|
if (
|
|
not main_model.supports_reasoning
|
|
or grant.reasoning_effort not in main_model.allowed_reasoning_efforts
|
|
):
|
|
raise EvoRuntimeError("MODEL_CAPABILITY_UNAVAILABLE")
|
|
result: dict[str, tuple[ResolvedRoute, ...]] = {}
|
|
result["main_agent"] = tuple(
|
|
self._resolve_route(
|
|
config, route, "main_agent", model_options=model_options or {}
|
|
)
|
|
for route in main_refs
|
|
)
|
|
for purpose in ("tool_selector", "deepagents_summarizer"):
|
|
selector_id = config.purpose_selector_ids.get(purpose)
|
|
refs = (
|
|
[await self._select_concrete(config, selector_id)]
|
|
if selector_id is not None
|
|
else main_refs
|
|
)
|
|
for ref in refs:
|
|
self._require_route_access(config, ref, grant)
|
|
result[purpose] = tuple(
|
|
self._resolve_route(
|
|
config,
|
|
route,
|
|
purpose,
|
|
model_options=model_options or {} if selector_id is None else {},
|
|
)
|
|
for route in refs
|
|
)
|
|
if grant.title_policy == "best_effort":
|
|
title_ref = await self._select_concrete(config, config.title_selector_id)
|
|
self._require_route_access(config, title_ref, grant)
|
|
result["title"] = (self._resolve_route(config, title_ref, "title"),)
|
|
return result
|
|
|
|
@classmethod
|
|
def _require_route_access(
|
|
cls,
|
|
config: EvoModelConfig,
|
|
route: RouteRef,
|
|
grant: RoutePreparationGrant,
|
|
) -> None:
|
|
selector = config.route_selectors[route.selector_id]
|
|
model = config.route_model(route)
|
|
alias_config = config.aliases.get(selector.alias)
|
|
if not cls._access_allowed(
|
|
model, grant.plan, grant.roles, grant.subject_id
|
|
) or (
|
|
alias_config is not None
|
|
and not cls._access_policy_allowed(
|
|
alias_config.access, grant.roles, grant.subject_id
|
|
)
|
|
):
|
|
raise EvoRuntimeError("MODEL_ACCESS_DENIED")
|
|
|
|
async def _select_concrete(
|
|
self, config: EvoModelConfig, selector_id: str
|
|
) -> RouteRef:
|
|
selector = config.route_selectors[selector_id]
|
|
routes = list(config.concrete_routes(selector_id))
|
|
closed = [
|
|
route
|
|
for route in routes
|
|
if self._route_health.status(route.key()) == "closed"
|
|
]
|
|
candidates = closed or [
|
|
route
|
|
for route in routes
|
|
if self._route_health.status(route.key()) == "half_open"
|
|
]
|
|
if not candidates:
|
|
raise EvoRuntimeError("MODEL_ROUTE_UNAVAILABLE")
|
|
if selector.endpoint_pool is None or len(candidates) == 1:
|
|
return candidates[0]
|
|
pool = config.endpoint_pools[selector.endpoint_pool]
|
|
candidate_names = {route.endpoint for route in candidates}
|
|
weighted = [
|
|
(member.name, member.weight)
|
|
for member in pool.endpoints
|
|
if member.name in candidate_names
|
|
]
|
|
chosen = await self._pool.choose(pool.pool_id, weighted)
|
|
return next(route for route in candidates if route.endpoint == chosen)
|
|
|
|
def _resolve_route(
|
|
self,
|
|
config: EvoModelConfig,
|
|
route: RouteRef,
|
|
purpose: str,
|
|
*,
|
|
model_options: Mapping[str, Any] | None = None,
|
|
) -> ResolvedRoute:
|
|
provider = config.providers[route.provider]
|
|
endpoint = provider.endpoints[route.endpoint]
|
|
model = provider.models[route.model]
|
|
semantics_key = self.identity_key_ring.derive(
|
|
config.config_identity_key_id, _ROUTE_SEMANTICS_INFO
|
|
)
|
|
fingerprint_key = self.identity_key_ring.derive(
|
|
config.config_identity_key_id, _ROUTE_FINGERPRINT_INFO
|
|
)
|
|
endpoint_key = self.identity_key_ring.derive(
|
|
config.config_identity_key_id, _ENDPOINT_FINGERPRINT_INFO
|
|
)
|
|
semantics_hash = route_semantics_hash(config, route, semantics_key)
|
|
auth = resolve_secret(endpoint.auth, secret_resolver=self.secret_resolver)
|
|
secret_fingerprints = {endpoint.auth.ref: auth.runtime_fingerprint}
|
|
headers = dict(endpoint.headers)
|
|
for header, reference in endpoint.header_refs.items():
|
|
resolved = resolve_secret(reference, secret_resolver=self.secret_resolver)
|
|
headers[header] = resolved.value
|
|
secret_fingerprints[reference.ref] = resolved.runtime_fingerprint
|
|
selector = config.route_selectors[route.selector_id]
|
|
if config.schema_version == 3:
|
|
params = dict(config.purpose_defaults.get(purpose, {}))
|
|
params.update(
|
|
project_user_options_for_purpose(
|
|
values=model.params,
|
|
user_options=model.user_options,
|
|
purpose=purpose,
|
|
)
|
|
)
|
|
for name, option in model.user_options.items():
|
|
if "default" in option and purpose in set(
|
|
option.get("applies_to") or ("main_agent",)
|
|
):
|
|
params.setdefault(name, option["default"])
|
|
params.update(
|
|
project_user_options_for_purpose(
|
|
values=selector.alias_defaults,
|
|
user_options=model.user_options,
|
|
purpose=purpose,
|
|
)
|
|
)
|
|
params.update(model.purpose_overrides.get(purpose, {}))
|
|
registration = get_adapter_registry().get(
|
|
provider.adapter_id, provider.adapter_revision
|
|
)
|
|
options = self._validated_user_options(
|
|
model,
|
|
registration,
|
|
project_user_options_for_purpose(
|
|
values=model_options or {},
|
|
user_options=model.user_options,
|
|
purpose=purpose,
|
|
),
|
|
purpose,
|
|
)
|
|
params.update(options)
|
|
registration.validate_parameters(params, path=f"invocation.{purpose}")
|
|
if registration.lifecycle == "blocked":
|
|
raise EvoRuntimeError("MODEL_ADAPTER_BLOCKED")
|
|
invocation_identity = invocation_fingerprint(
|
|
config, route, purpose, params, fingerprint_key
|
|
)
|
|
else:
|
|
params = {
|
|
**config.runtime_defaults,
|
|
**provider.params,
|
|
**endpoint.params,
|
|
**model.params,
|
|
}
|
|
registration = None
|
|
invocation_identity = route_fingerprint(config, route, fingerprint_key)
|
|
supports_tools = bool(model.capabilities.get("tools", False))
|
|
effective_tool_transport = (
|
|
"native" if supports_tools else "disabled"
|
|
)
|
|
identity = RouteIdentity(
|
|
config_revision=config.config_revision,
|
|
config_identity_key_id=config.config_identity_key_id,
|
|
purpose=purpose,
|
|
route_selector_id=selector.identity_selector_id or route.selector_id,
|
|
route_fingerprint=invocation_identity,
|
|
provider_id=route.provider,
|
|
endpoint_name=route.endpoint,
|
|
model_id=model.model_id,
|
|
protocol=provider.adapter_id or provider.protocol,
|
|
api_mode=route.api_mode,
|
|
tool_call_transport=effective_tool_transport,
|
|
route_semantics_hash=semantics_hash,
|
|
billing_sku=model.billing_sku,
|
|
pricing_revision=model.quote.pricing_revision,
|
|
quote_id=model.quote.quote_id,
|
|
)
|
|
return ResolvedRoute(
|
|
route,
|
|
identity,
|
|
model.quote,
|
|
model.context_window,
|
|
model.max_output_tokens,
|
|
model.reasoning_mode,
|
|
model.reasoning_enabled_params,
|
|
model.reasoning_disabled_params,
|
|
endpoint.base_url,
|
|
params,
|
|
auth.value,
|
|
headers,
|
|
secret_fingerprints,
|
|
registration,
|
|
int(provider.connection_defaults.get("max_inflight_requests", 16)),
|
|
int(
|
|
model.max_inflight_requests
|
|
or provider.connection_defaults.get("max_inflight_requests", 16)
|
|
),
|
|
int(provider.connection_defaults.get("queue_timeout_seconds", 5)),
|
|
f"provider:{route.provider}:{endpoint_fingerprint(config, route, endpoint_key)}",
|
|
f"model:{route.key()}:{semantics_hash}",
|
|
int(provider.connection_defaults.get("attempt_timeout_seconds", 600)),
|
|
supports_tools=supports_tools,
|
|
)
|
|
|
|
@staticmethod
|
|
def _validated_user_options(
|
|
model: Any,
|
|
registration: AdapterRegistration,
|
|
supplied: Mapping[str, Any],
|
|
purpose: str,
|
|
) -> Mapping[str, Any]:
|
|
if {"reasoning", "reasoning_effort"} & set(supplied):
|
|
raise EvoRuntimeError("AGENT_INPUT_MISMATCH")
|
|
validated = validate_user_model_options(
|
|
supplied=supplied,
|
|
user_options=model.user_options,
|
|
supports_reasoning=model.supports_reasoning,
|
|
reasoning_mode=model.reasoning_mode,
|
|
allowed_reasoning_efforts=model.allowed_reasoning_efforts,
|
|
default_reasoning_effort=str(
|
|
model.reasoning_enabled_params.get("reasoning") or ""
|
|
),
|
|
parameter_constraints=model.parameter_constraints,
|
|
purpose=purpose,
|
|
allow_reasoning=False,
|
|
)
|
|
for name, value in validated.items():
|
|
registration.all_parameter_schema[name].validate(
|
|
value, f"model_options.{name}"
|
|
)
|
|
return validated
|
|
|
|
def _build_call_bounds(
|
|
self,
|
|
purpose_routes: Mapping[str, tuple[ResolvedRoute, ...]],
|
|
) -> Mapping[str, tuple[RouteCallBound, ...]]:
|
|
result: dict[str, tuple[RouteCallBound, ...]] = {}
|
|
for purpose, routes in purpose_routes.items():
|
|
bounds = []
|
|
for route in routes:
|
|
override = route.params.get("output_token_limit")
|
|
output_limit = (
|
|
route.max_output_tokens if override is None else int(override)
|
|
)
|
|
if output_limit > route.max_output_tokens:
|
|
raise EvoRuntimeError("MODEL_OUTPUT_LIMIT_EXCEEDED")
|
|
margin = _PROTOCOL_MARGIN_TOKENS.get(
|
|
(route.identity.protocol, route.identity.api_mode)
|
|
)
|
|
if margin is None:
|
|
raise EvoRuntimeError("TOKEN_BOUND_UNAVAILABLE")
|
|
billable_cap = route.context_window - output_limit
|
|
payload_cap = billable_cap - margin
|
|
if payload_cap <= 0:
|
|
raise EvoRuntimeError("MODEL_CONTEXT_WINDOW_EXCEEDED")
|
|
reserve_rate = max(
|
|
route.quote.input_microunits_per_million,
|
|
route.quote.cached_input_microunits_per_million,
|
|
)
|
|
reserve = _checked_ceil_cost(
|
|
billable_cap,
|
|
reserve_rate,
|
|
output_limit,
|
|
route.quote.output_microunits_per_million,
|
|
route.quote.unit_scale,
|
|
)
|
|
bounds.append(
|
|
RouteCallBound(
|
|
route.identity,
|
|
output_limit,
|
|
payload_cap,
|
|
billable_cap,
|
|
margin,
|
|
reserve,
|
|
)
|
|
)
|
|
result[purpose] = tuple(bounds)
|
|
return result
|
|
|
|
@staticmethod
|
|
def _run_reserve(
|
|
attempt_limits: Mapping[str, int],
|
|
bounds: Mapping[str, tuple[RouteCallBound, ...]],
|
|
) -> int:
|
|
"""Do not turn model-call estimates into a runtime admission limit.
|
|
|
|
Agent progress is bounded by LangGraph's per-run ``recursion_limit``.
|
|
The legacy purpose counts remain in the signed quote for wire
|
|
compatibility and reporting, but a normal model/tool loop must not be
|
|
stopped because it outgrows that estimate.
|
|
"""
|
|
del attempt_limits, bounds
|
|
return 0
|
|
|
|
def _build_model_set(
|
|
self, snapshot: ModelRuntimeSnapshot, user_reasoning_effort: str
|
|
) -> AgentModelSet:
|
|
del user_reasoning_effort
|
|
|
|
def models(purpose: str) -> tuple[Any, ...]:
|
|
return tuple(
|
|
self._build_model(route, purpose, snapshot)
|
|
for route in snapshot.purpose_routes[purpose]
|
|
)
|
|
|
|
main = models("main_agent")
|
|
# A fallback chain is a route choice, not a workflow-call budget.
|
|
# Do not duplicate the primary model from max_attempts_per_run: that
|
|
# field is legacy quote metadata and must not cap normal agent turns.
|
|
retry_models = main[1:]
|
|
return AgentModelSet(
|
|
main_agent=main[0],
|
|
tool_selector=models("tool_selector")[0],
|
|
deepagents_summarizer=models("deepagents_summarizer")[0],
|
|
title=models("title")[0] if "title" in snapshot.purpose_routes else None,
|
|
main_fallbacks=retry_models,
|
|
route_health=self._route_health,
|
|
capacity=self._capacity,
|
|
)
|
|
|
|
def _build_model(
|
|
self,
|
|
route: ResolvedRoute,
|
|
purpose: str,
|
|
snapshot: ModelRuntimeSnapshot,
|
|
) -> Any:
|
|
bound = next(
|
|
item
|
|
for item in snapshot.purpose_route_call_bounds[purpose]
|
|
if item.route_identity == route.identity
|
|
)
|
|
if not route.runtime_provider or route.invocation_plan is None:
|
|
raise EvoRuntimeError("MODEL_ADAPTER_COMPILE_FAILED")
|
|
model = self.model_factory(
|
|
model=route.identity.model_id,
|
|
provider=route.runtime_provider,
|
|
**self._model_factory_kwargs(route),
|
|
)
|
|
if isinstance(model, ModelFactoryResult):
|
|
owner = _construction_owner.get()
|
|
if owner is None and model.owned_clients:
|
|
raise EvoRuntimeError("MODEL_RESOURCE_OWNER_MISSING")
|
|
if owner is not None:
|
|
for client in model.owned_clients:
|
|
owner._owned_clients[id(client)] = client
|
|
model = model.model
|
|
return self._attach_route_metadata(model, route, purpose, bound)
|
|
|
|
@staticmethod
|
|
def _model_factory_kwargs(route: ResolvedRoute) -> dict[str, Any]:
|
|
plan = route.invocation_plan
|
|
if plan is None:
|
|
raise EvoRuntimeError("MODEL_ADAPTER_COMPILE_FAILED")
|
|
params = plan.model_kwargs()
|
|
params.update(
|
|
{
|
|
"api_key": route.api_key,
|
|
"base_url": route.base_url,
|
|
"max_retries": 0,
|
|
"timeout": route.attempt_timeout_seconds,
|
|
}
|
|
)
|
|
if route.default_headers:
|
|
params["default_headers"] = dict(route.default_headers)
|
|
return params
|
|
|
|
def _compile_purpose_routes(
|
|
self,
|
|
purpose_routes: Mapping[str, tuple[ResolvedRoute, ...]],
|
|
bounds: Mapping[str, tuple[RouteCallBound, ...]],
|
|
user_reasoning_effort: str,
|
|
) -> Mapping[str, tuple[ResolvedRoute, ...]]:
|
|
"""Compile each frozen purpose route exactly once during Prepare."""
|
|
|
|
compiled: dict[str, tuple[ResolvedRoute, ...]] = {}
|
|
for purpose, routes in purpose_routes.items():
|
|
by_identity = {item.route_identity: item for item in bounds[purpose]}
|
|
compiled[purpose] = tuple(
|
|
self._compile_route(
|
|
route,
|
|
purpose,
|
|
by_identity[route.identity],
|
|
user_reasoning_effort,
|
|
)
|
|
for route in routes
|
|
)
|
|
return compiled
|
|
|
|
@staticmethod
|
|
def _compile_route(
|
|
route: ResolvedRoute,
|
|
purpose: str,
|
|
bound: RouteCallBound,
|
|
user_reasoning_effort: str,
|
|
) -> ResolvedRoute:
|
|
params = dict(route.params)
|
|
effort = user_reasoning_effort if purpose == "main_agent" else "disabled"
|
|
if route.adapter is not None:
|
|
if effort != "disabled":
|
|
params["reasoning_effort"] = effort
|
|
else:
|
|
params["thinking_enabled"] = False
|
|
compiled = dict(
|
|
route.adapter.compile_runtime_parameters(
|
|
route.identity.api_mode,
|
|
params,
|
|
bound.max_output_tokens,
|
|
provider_model_id=route.identity.model_id,
|
|
)
|
|
)
|
|
runtime_provider = (
|
|
"google_interactions"
|
|
if route.adapter.adapter_id == "google-gemini"
|
|
and route.identity.api_mode == "interactions"
|
|
else route.adapter.runtime_provider
|
|
)
|
|
plan = compile_invocation_plan(
|
|
api_mode=route.identity.api_mode,
|
|
declared_tool_call_transport=route.identity.tool_call_transport,
|
|
supports_tools=route.supports_tools,
|
|
purpose=purpose,
|
|
output_token_limit=bound.max_output_tokens,
|
|
reasoning_effort=effort,
|
|
runtime_provider=runtime_provider,
|
|
sdk_params=compiled,
|
|
)
|
|
return replace(
|
|
route,
|
|
runtime_provider=runtime_provider,
|
|
invocation_plan=plan,
|
|
)
|
|
|
|
params.pop("output_token_limit", None)
|
|
params.update(
|
|
{
|
|
"max_tokens": bound.max_output_tokens,
|
|
"use_responses_api": route.identity.api_mode == "responses",
|
|
}
|
|
)
|
|
if route.default_headers:
|
|
params["default_headers"] = dict(route.default_headers)
|
|
if effort != "disabled" and route.reasoning_mode == "effort":
|
|
params["reasoning_effort"] = effort
|
|
elif route.reasoning_mode == "boolean":
|
|
params.pop("reasoning_effort", None)
|
|
params = _merge_params(
|
|
params,
|
|
route.reasoning_enabled_params
|
|
if effort != "disabled"
|
|
else route.reasoning_disabled_params,
|
|
)
|
|
else:
|
|
params.pop("reasoning_effort", None)
|
|
plan = compile_invocation_plan(
|
|
api_mode=route.identity.api_mode,
|
|
declared_tool_call_transport=route.identity.tool_call_transport,
|
|
supports_tools=route.supports_tools,
|
|
purpose=purpose,
|
|
output_token_limit=bound.max_output_tokens,
|
|
reasoning_effort=effort,
|
|
runtime_provider=route.identity.protocol,
|
|
sdk_params=params,
|
|
)
|
|
return replace(
|
|
route,
|
|
runtime_provider=route.identity.protocol,
|
|
invocation_plan=plan,
|
|
)
|
|
|
|
@staticmethod
|
|
def _attach_route_metadata(
|
|
model: Any, route: ResolvedRoute, purpose: str, bound: RouteCallBound
|
|
) -> Any:
|
|
metadata = {
|
|
"route_key": route.identity.route_key,
|
|
"route_provider": route.identity.provider_id,
|
|
"route_model": route.identity.model_id,
|
|
"route_endpoint": route.identity.endpoint_name,
|
|
"route_api_mode": route.invocation_plan.api_mode
|
|
if route.invocation_plan
|
|
else route.identity.api_mode,
|
|
"route_tool_call_transport": route.invocation_plan.tool_call_transport
|
|
if route.invocation_plan
|
|
else route.identity.tool_call_transport,
|
|
"route_supports_tools": route.supports_tools,
|
|
"route_invocation_plan_hash": route.invocation_plan.plan_hash
|
|
if route.invocation_plan
|
|
else "",
|
|
"route_output_token_parameter": route.invocation_plan.output_token_parameter
|
|
if route.invocation_plan
|
|
else "",
|
|
"route_adapter_id": route.adapter.adapter_id if route.adapter else "",
|
|
"route_adapter_revision": route.adapter.adapter_revision
|
|
if route.adapter
|
|
else "",
|
|
"route_config_generation": route.identity.config_revision,
|
|
"runtime_purpose": purpose,
|
|
"runtime_output_limit": bound.max_output_tokens,
|
|
"capacity_provider_key": route.provider_capacity_key,
|
|
"capacity_model_key": route.model_capacity_key,
|
|
"provider_max_inflight": route.provider_max_inflight,
|
|
"model_max_inflight": route.model_max_inflight,
|
|
"capacity_queue_timeout_seconds": route.queue_timeout_seconds,
|
|
"attempt_timeout_seconds": route.attempt_timeout_seconds,
|
|
}
|
|
if hasattr(model, "model_copy"):
|
|
model = model.model_copy(
|
|
update={
|
|
"metadata": {**(getattr(model, "metadata", None) or {}), **metadata}
|
|
}
|
|
)
|
|
return model
|
|
|
|
def _validate_preparation_echo(
|
|
self, grant: RoutePreparationGrant, agent_input: AgentInputV3
|
|
) -> None:
|
|
if (
|
|
agent_input.checkpoint_thread_id != grant.checkpoint_thread_id
|
|
or self.quote_authority.agent_input_digest(agent_input.projection())
|
|
!= grant.gateway_input_digest
|
|
or grant.turn_fencing_token < 1
|
|
or grant.reasoning_effort
|
|
not in {"disabled", "low", "medium", "high", "max"}
|
|
):
|
|
raise EvoRuntimeError("AGENT_INPUT_MISMATCH")
|
|
|
|
@staticmethod
|
|
def _validate_admission_echo(
|
|
quote: PreparedRunQuote, admission: AdmissionGrant
|
|
) -> None:
|
|
names = (
|
|
"preparation_id",
|
|
"request_id",
|
|
"turn_id",
|
|
"thread_id",
|
|
"subject_id",
|
|
"requested_model_ref",
|
|
"plan",
|
|
"roles",
|
|
"requires_vision",
|
|
"reasoning_effort",
|
|
"title_policy",
|
|
"gateway_input_digest",
|
|
"prepared_snapshot_digest",
|
|
"prepared_input_digest",
|
|
"config_revision",
|
|
"catalog_revision",
|
|
"purpose_attempt_limits",
|
|
"total_max_attempts",
|
|
"checkpoint_snapshot_id",
|
|
"tool_registry_snapshot_id",
|
|
"turn_fencing_token",
|
|
"provider_run_reserve_microunits",
|
|
)
|
|
if any(getattr(quote, name) != getattr(admission, name) for name in names):
|
|
raise EvoRuntimeError("ADMISSION_INVALID")
|
|
if (
|
|
admission.billing_fencing_token < 1
|
|
or not admission.admission_snapshot_id
|
|
or not admission.admission_id
|
|
or not admission.hold_id
|
|
):
|
|
raise EvoRuntimeError("ADMISSION_INVALID")
|
|
|
|
def _validate_still_fresh(self, handle: _PreparedHandle) -> None:
|
|
config = self.store.load_revision(handle.snapshot.config_revision)
|
|
if self._tool_registry_payload(handle.host) != handle.tool_registry_payload:
|
|
raise EvoRuntimeError("TOOL_REGISTRY_STALE")
|
|
for routes in handle.snapshot.purpose_routes.values():
|
|
for route in routes:
|
|
endpoint = config.providers[route.ref.provider].endpoints[
|
|
route.ref.endpoint
|
|
]
|
|
for reference in (endpoint.auth, *endpoint.header_refs.values()):
|
|
resolved = resolve_secret(
|
|
reference, secret_resolver=self.secret_resolver
|
|
)
|
|
if (
|
|
handle.secret_fingerprints.get(reference.ref)
|
|
!= resolved.runtime_fingerprint
|
|
):
|
|
raise EvoRuntimeError("PREPARATION_STALE")
|
|
|
|
@staticmethod
|
|
def _tool_registry_payload(host: WebHostContext) -> Mapping[str, Any]:
|
|
registry = host.tool_registry
|
|
revision = host.tool_registry_revision
|
|
if host.tool_registry_provider is not None:
|
|
registry, revision = host.tool_registry_provider()
|
|
tools = []
|
|
for tool in registry:
|
|
if isinstance(tool, Mapping):
|
|
tools.append(
|
|
{
|
|
"name": str(tool.get("name") or ""),
|
|
"description": str(tool.get("description") or ""),
|
|
"schema": dict(tool.get("schema") or {}),
|
|
}
|
|
)
|
|
continue
|
|
schema = getattr(tool, "args_schema", None)
|
|
if hasattr(schema, "model_json_schema"):
|
|
schema = schema.model_json_schema()
|
|
tools.append(
|
|
{
|
|
"name": str(
|
|
getattr(
|
|
tool, "name", getattr(tool, "__name__", type(tool).__name__)
|
|
)
|
|
),
|
|
"description": str(getattr(tool, "description", "")),
|
|
"schema": schema or {},
|
|
}
|
|
)
|
|
return {"revision": revision, "tools": tools}
|
|
|
|
def _expire_prepared_locked(self) -> None:
|
|
current = now_ms()
|
|
for handle in tuple(self._prepared.values()):
|
|
if handle.state == "PREPARED" and handle.quote.expires_at < current:
|
|
handle.state = "EXPIRED"
|
|
self._retire_prepared_locked(handle)
|
|
for key, tombstone in tuple(self._prepared_tombstones.items()):
|
|
if tombstone.retained_until < current:
|
|
self._prepared_tombstones.pop(key, None)
|
|
|
|
def _prepared_replay_locked(
|
|
self, request_key: tuple[str, str, str], request_digest: str
|
|
) -> PreparedRunQuote | None:
|
|
tombstone = self._prepared_tombstones.get(request_key)
|
|
if tombstone is not None:
|
|
if tombstone.request_digest != request_digest:
|
|
raise EvoRuntimeError("PREPARATION_CONFLICT")
|
|
raise EvoRuntimeError("PREPARATION_STALE")
|
|
preparation_id = self._prepared_requests.get(request_key)
|
|
if preparation_id is None:
|
|
return None
|
|
handle = self._prepared.get(preparation_id)
|
|
if handle is None:
|
|
self._prepared_requests.pop(request_key, None)
|
|
return None
|
|
if handle.state in {"CANCELLED", "EXPIRED"}:
|
|
self._retire_prepared_locked(handle)
|
|
if handle.request_digest != request_digest:
|
|
raise EvoRuntimeError("PREPARATION_CONFLICT")
|
|
raise EvoRuntimeError("PREPARATION_STALE")
|
|
if handle.request_digest != request_digest:
|
|
raise EvoRuntimeError("PREPARATION_CONFLICT")
|
|
return handle.quote
|
|
|
|
def _retire_prepared_locked(self, handle: _PreparedHandle) -> None:
|
|
request_key = (
|
|
handle.grant.subject_id,
|
|
handle.grant.request_id,
|
|
handle.grant.turn_id,
|
|
)
|
|
self._prepared_tombstones[request_key] = _PreparedTombstone(
|
|
handle.request_digest,
|
|
handle.state,
|
|
now_ms() + 24 * 60 * 60 * 1000,
|
|
)
|
|
self._prepared_requests.pop(request_key, None)
|
|
self._prepared.pop(handle.quote.preparation_id, None)
|
|
|
|
@staticmethod
|
|
def _preparation_request_digest(
|
|
grant: RoutePreparationGrant,
|
|
agent_input: AgentInputV3,
|
|
tool_registry_payload: Mapping[str, Any],
|
|
) -> str:
|
|
grant_payload = grant.unsigned_payload()
|
|
for ephemeral in ("grant_id", "issued_at", "expires_at", "key_id"):
|
|
grant_payload.pop(ephemeral, None)
|
|
payload = {
|
|
"grant": grant_payload,
|
|
"agent_input": agent_input.projection(),
|
|
"tool_registry": tool_registry_payload,
|
|
}
|
|
return hashlib.sha256(canonical_json_v1(payload)).hexdigest()
|
|
|
|
@staticmethod
|
|
def _access_allowed(
|
|
model: Any, plan: str, roles: tuple[str, ...], subject_id: str = ""
|
|
) -> bool:
|
|
if model.allowed_plans and plan not in model.allowed_plans:
|
|
return False
|
|
if model.allowed_roles and not bool(set(roles) & set(model.allowed_roles)):
|
|
return False
|
|
return EvoModelRuntime._access_policy_allowed(
|
|
getattr(model, "access", {}), roles, subject_id
|
|
)
|
|
|
|
@staticmethod
|
|
def _access_policy_allowed(
|
|
policy: Mapping[str, Any], roles: tuple[str, ...], subject_id: str
|
|
) -> bool:
|
|
if not policy or policy.get("visibility") == "authenticated":
|
|
return True
|
|
if policy.get("visibility") == "role_based":
|
|
role_set = set(roles)
|
|
group_set = {
|
|
role.removeprefix("group:")
|
|
for role in role_set
|
|
if role.startswith("group:")
|
|
}
|
|
return bool(
|
|
role_set & set(policy.get("roles") or ())
|
|
or group_set & set(policy.get("groups") or ())
|
|
)
|
|
if policy.get("visibility") == "private":
|
|
return subject_id in set(policy.get("users") or ())
|
|
return False
|
|
|
|
@staticmethod
|
|
def _default_model_factory(**kwargs: Any) -> Any:
|
|
if kwargs.get("provider") == "google_interactions":
|
|
from .gemini_interactions import create_gemini_interactions_model
|
|
|
|
return create_gemini_interactions_model(**kwargs)
|
|
from .models import _OPENAI_ROUTED_PROVIDERS, get_chat_model
|
|
|
|
# LangChain caches default OpenAI HTTP clients process-wide. A run must
|
|
# never close that cache underneath another run's active stream.
|
|
if kwargs.get("provider") == "openai" or kwargs.get("provider") in _OPENAI_ROUTED_PROVIDERS:
|
|
import httpx
|
|
|
|
# Supplying owned clients disables LangChain's automatic usage opt-in.
|
|
# Keep explicit caller configuration and borrowed transports intact.
|
|
kwargs.setdefault("stream_usage", True)
|
|
if "http_client" not in kwargs:
|
|
kwargs["http_client"] = httpx.Client()
|
|
owner = _construction_owner.get()
|
|
if owner is not None:
|
|
owner._owned_clients[id(kwargs["http_client"])] = kwargs["http_client"]
|
|
if "http_async_client" not in kwargs:
|
|
kwargs["http_async_client"] = httpx.AsyncClient()
|
|
owner = _construction_owner.get()
|
|
if owner is not None:
|
|
owner._owned_clients[id(kwargs["http_async_client"])] = kwargs["http_async_client"]
|
|
return get_chat_model(**kwargs)
|
|
|
|
@staticmethod
|
|
def _default_agent_factory(
|
|
snapshot: ModelRuntimeSnapshot, host: WebHostContext, model_set: AgentModelSet
|
|
) -> Any:
|
|
from ..web_runtime import create_web_agent
|
|
|
|
return create_web_agent(snapshot=snapshot, host=host, model_set=model_set)
|
|
|
|
|
|
class _EvoWebRun:
|
|
def __init__(
|
|
self,
|
|
*,
|
|
runtime: EvoModelRuntime,
|
|
admission: AdmissionGrant,
|
|
prepared: _PreparedHandle,
|
|
agent: Any,
|
|
model_set: AgentModelSet,
|
|
) -> None:
|
|
self._runtime = runtime
|
|
self._admission = admission
|
|
self._prepared = prepared
|
|
self._snapshot = prepared.snapshot
|
|
self._input = prepared.input
|
|
self._host = prepared.host
|
|
self._agent = agent
|
|
self._model_set = model_set
|
|
self._run_id = prepared.quote.execution_id or str(uuid.uuid4())
|
|
self._state = "ACTIVE"
|
|
self._sequence = 0
|
|
self._journal: list[EvoRuntimeEvent] = []
|
|
self._journal_bytes = 0
|
|
self._condition = asyncio.Condition()
|
|
self._event_commit_lock = asyncio.Lock()
|
|
self._budget_lock = asyncio.Lock()
|
|
self._terminal_lock = asyncio.Lock()
|
|
self._agent_task: asyncio.Task[None] | None = None
|
|
self._checkpoint_writes_stopped = True
|
|
from ..stream.stop import StreamStop
|
|
self._stream_stop = StreamStop() if runtime.host_registry is not None else None
|
|
self._event_iterator = None
|
|
self._stop_task = None
|
|
self._stop_errors = []
|
|
self._stream_wakeup_task: asyncio.Task[None] | None = None
|
|
self._attempt_counts: dict[str, int] = defaultdict(int)
|
|
self._remaining = admission.provider_run_reserve_microunits
|
|
self._callback_attempts: dict[
|
|
str, tuple[ModelAttemptEvent, int, PricingQuote, ResolvedRoute]
|
|
] = {}
|
|
self._last_model_failure: tuple[str, Mapping[str, Any]] | None = None
|
|
self._force_context_repair_pending = bool(
|
|
self._input.metadata.get("force_context_repair", False)
|
|
)
|
|
self._attempt_callback = _RuntimeAttemptCallback(self)
|
|
self._terminal_event: EvoRuntimeEvent | None = None
|
|
self._terminal_task: asyncio.Task[EvoRuntimeEvent] | None = None
|
|
self._terminal_error_code: str | None = None
|
|
self._terminal_phase = "cleanup"
|
|
self._cancel_task: asyncio.Task[None] | None = None
|
|
self._cleanup_task: asyncio.Task[None] | None = None
|
|
self._resource_close_tasks: dict[int, asyncio.Task[None]] = {}
|
|
self._owned_clients: dict[int, Any] = {}
|
|
self._cleanup_failures: dict[int, list[BaseException]] = defaultdict(list)
|
|
|
|
async def _close_owned_client(self, key: int, client: Any) -> None:
|
|
# A pending attempt owns the resource until it really returns. In
|
|
# particular, cancelling an observer cannot stop a synchronous thread.
|
|
for attempt in range(3):
|
|
try:
|
|
close = getattr(client, "aclose", None) or getattr(client, "close")
|
|
if asyncio.iscoroutinefunction(close):
|
|
result = close()
|
|
else:
|
|
result = await asyncio.to_thread(close)
|
|
if hasattr(result, "__await__"):
|
|
await result
|
|
except BaseException as exc:
|
|
self._cleanup_failures[key].append(exc)
|
|
if isinstance(exc, asyncio.CancelledError):
|
|
return
|
|
else:
|
|
self._owned_clients.pop(key, None)
|
|
return
|
|
if attempt < 2:
|
|
await asyncio.sleep(0)
|
|
|
|
async def _close_owned_clients(self) -> None:
|
|
self._resource_close_tasks = {
|
|
key: asyncio.create_task(
|
|
self._close_owned_client(key, client),
|
|
name=f"evo-resource-close-{self._run_id}-{key}",
|
|
)
|
|
for key, client in tuple(self._owned_clients.items())
|
|
}
|
|
if self._resource_close_tasks:
|
|
await asyncio.gather(*self._resource_close_tasks.values())
|
|
|
|
|
|
async def _ensure_cleanup(self, *, timeout: float | None = 0.25) -> None:
|
|
if self._cleanup_task is None:
|
|
self._cleanup_task = asyncio.create_task(
|
|
self._close_owned_clients(), name=f"evo-cleanup-{self._run_id}"
|
|
)
|
|
# asyncio.wait bounds observation without cancelling the run-owned
|
|
# coordinator or any close operation, including cancellation-resistant
|
|
# async clients. Later observers can confirm eventual completion.
|
|
await asyncio.wait({self._cleanup_task}, timeout=timeout)
|
|
if not self._cleanup_task.done():
|
|
raise TimeoutError("Run resource close remains unconfirmed")
|
|
self._cleanup_task.result()
|
|
if self._owned_clients:
|
|
errors = []
|
|
for key in self._owned_clients:
|
|
task = self._resource_close_tasks.get(key)
|
|
failures = self._cleanup_failures.get(key, [])
|
|
errors.append(
|
|
failures[-1] if task is not None and task.done() and failures
|
|
else TimeoutError("Run resource close remains unconfirmed")
|
|
)
|
|
raise BaseExceptionGroup("Run resource cleanup unconfirmed", errors)
|
|
|
|
def _start(self) -> None:
|
|
if self._agent_task is None and self._terminal_event is None:
|
|
self._agent_task = asyncio.create_task(
|
|
self._run_with_timeout(), name=f"evo-run-{self._run_id}"
|
|
)
|
|
self._agent_task.add_done_callback(self._wake_stream_waiters)
|
|
|
|
@property
|
|
def run_id(self) -> str:
|
|
return self._run_id
|
|
|
|
async def stream(
|
|
self, after_sequence: int | None = None
|
|
) -> AsyncIterator[EvoRuntimeEvent]:
|
|
async with self._condition:
|
|
if after_sequence is not None and (
|
|
isinstance(after_sequence, bool) or not isinstance(after_sequence, int)
|
|
):
|
|
raise EvoRuntimeError("EVENT_CURSOR_INVALID")
|
|
cursor = 0 if after_sequence is None else after_sequence
|
|
if cursor < 0 or cursor > self._sequence:
|
|
raise EvoRuntimeError("EVENT_CURSOR_INVALID")
|
|
|
|
async for event in self._replay(cursor):
|
|
yield event
|
|
|
|
def _wake_stream_waiters(self, _task: asyncio.Task[None]) -> None:
|
|
self._stream_wakeup_task = asyncio.create_task(self._notify_stream_waiters())
|
|
|
|
async def _notify_stream_waiters(self) -> None:
|
|
async with self._condition:
|
|
self._condition.notify_all()
|
|
|
|
async def cancel(self, reason: str, *, owner_epoch: int | None = None,
|
|
boot_id: str | None = None) -> str:
|
|
if self._runtime.host_registry is not None:
|
|
return await self._runtime.cancel(
|
|
self.run_id, reason=reason, owner_epoch=owner_epoch, boot_id=boot_id)
|
|
self._request_cancel(reason)
|
|
return await self.wait_stopped(timeout=2.0)
|
|
|
|
def _request_cancel(self, reason: str) -> None:
|
|
if self._terminal_task is not None or self._terminal_event is not None:
|
|
return
|
|
if self._cancel_task is None:
|
|
self._state = "CANCELLING"
|
|
if self._agent_task is not None and not self._agent_task.done():
|
|
self._agent_task.cancel()
|
|
self._cancel_task = asyncio.create_task(
|
|
self._finish_cancel(reason), name=f"evo-stop-{self._run_id}"
|
|
)
|
|
|
|
async def wait_stopped(self, timeout: float | None = 2.0) -> str:
|
|
if self._terminal_event is not None:
|
|
return str(self._terminal_event.payload.get("outcome") or "terminal")
|
|
try:
|
|
async with asyncio.timeout(timeout):
|
|
task = self._cancel_task or self._agent_task
|
|
if task is not None:
|
|
await asyncio.shield(task)
|
|
if self._terminal_task is not None:
|
|
await asyncio.shield(self._terminal_task)
|
|
except Exception:
|
|
# Observation is not an ownership transfer or a stop proof.
|
|
return "unknown"
|
|
if self._terminal_event is None:
|
|
return "unknown"
|
|
return str(self._terminal_event.payload.get("outcome") or "terminal")
|
|
|
|
async def _finish_cancel(self, reason: str) -> None:
|
|
if self._agent_task is not None:
|
|
await asyncio.gather(self._agent_task, return_exceptions=True)
|
|
await self._terminal_locked(
|
|
"cancelled", error_code="RUN_CANCELLED", reason=reason
|
|
)
|
|
|
|
async def _run_with_timeout(self) -> None:
|
|
try:
|
|
async with asyncio.timeout(self._snapshot.active_run_timeout_seconds):
|
|
await self._run_agent()
|
|
except TimeoutError:
|
|
await self._terminal_locked("failed", error_code="RUN_TIMEOUT")
|
|
except asyncio.CancelledError:
|
|
return
|
|
|
|
async def _run_agent(self) -> None:
|
|
try:
|
|
await self._append_locked(
|
|
"run",
|
|
{
|
|
"kind": "run_started",
|
|
"outcome": "active",
|
|
"preparation_id": self._admission.preparation_id,
|
|
"admission_snapshot_id": self._admission.admission_snapshot_id,
|
|
"admission_id": self._admission.admission_id,
|
|
"hold_id": self._admission.hold_id,
|
|
"turn_fencing_token": self._admission.turn_fencing_token,
|
|
"billing_fencing_token": self._admission.billing_fencing_token,
|
|
},
|
|
)
|
|
from ..stream.events import (
|
|
_snapshot_has_pending_interrupt,
|
|
stream_agent_events,
|
|
)
|
|
|
|
self._checkpoint_writes_stopped = False
|
|
self._event_iterator = stream_agent_events(
|
|
self._agent,
|
|
self._input.message,
|
|
self._input.checkpoint_thread_id,
|
|
metadata=dict(self._input.metadata),
|
|
media=[
|
|
item.locator
|
|
if hasattr(item, "locator")
|
|
else str(item.get("locator") or "")
|
|
if isinstance(item, Mapping)
|
|
else str(item)
|
|
for item in self._input.media
|
|
],
|
|
callbacks=[self._attempt_callback],
|
|
configurable={
|
|
"turn_fencing_token": self._admission.turn_fencing_token,
|
|
"turn_lease_owner": self._admission.request_id,
|
|
"checkpoint_writer_strict_close": self._runtime.host_registry is not None,
|
|
**({"checkpoint_id": self._prepared.grant.predecessor_checkpoint_id,
|
|
"checkpoint_ns": ""}
|
|
if self._prepared.grant.predecessor_execution_id else {}),
|
|
},
|
|
error_mode="raise",
|
|
**({"stop_owner": self._stream_stop} if self._stream_stop else {}),
|
|
)
|
|
try:
|
|
async for source in self._event_iterator:
|
|
if str(source.get("type") or "") == "done":
|
|
continue
|
|
await self._append_locked("agent", dict(source))
|
|
finally:
|
|
if self._stop_task is None:
|
|
self._stop_task = asyncio.create_task(
|
|
self._stop_checkpoint_writer(), name=f"evo-writer-stop-{self._run_id}")
|
|
await asyncio.shield(self._stop_task)
|
|
try:
|
|
final_state = await self._agent.aget_state(
|
|
{"configurable": {"thread_id": self._input.checkpoint_thread_id}}
|
|
)
|
|
checkpoint = final_state.config["configurable"]
|
|
if not checkpoint.get("checkpoint_id"):
|
|
raise ValueError("Final checkpoint identity missing")
|
|
pending = []
|
|
for item in (
|
|
*getattr(final_state, "interrupts", ()),
|
|
*(item for task in getattr(final_state, "tasks", ())
|
|
for item in getattr(task, "interrupts", ())),
|
|
):
|
|
projected = {"id": item.id, "value": item.value}
|
|
if projected not in pending:
|
|
pending.append(projected)
|
|
checkpoint_details = json.loads(json.dumps({
|
|
"checkpoint_thread_id": checkpoint["thread_id"],
|
|
"checkpoint_id": checkpoint["checkpoint_id"],
|
|
"checkpoint_ns": checkpoint.get("checkpoint_ns", ""),
|
|
"pending_interrupts": pending,
|
|
}, allow_nan=False))
|
|
awaiting_input = _snapshot_has_pending_interrupt(final_state)
|
|
except Exception:
|
|
await self._terminal_locked(
|
|
"failed", error_code="FINAL_CHECKPOINT_READ_FAILED"
|
|
)
|
|
return
|
|
if awaiting_input:
|
|
await self._terminal_locked(
|
|
"awaiting_input", checkpoint_details=checkpoint_details
|
|
)
|
|
return
|
|
if self._admission.title_policy == "best_effort":
|
|
await self._run_title()
|
|
else:
|
|
await self._append_locked(
|
|
"title", {"kind": "skipped", "reason": "DISABLED"}
|
|
)
|
|
await self._terminal_locked("completed", checkpoint_details=checkpoint_details)
|
|
except asyncio.CancelledError:
|
|
raise
|
|
except Exception as exc:
|
|
logger.exception(
|
|
"evo run %s failed before model provider error was reported",
|
|
getattr(self._admission, "request_id", "?"),
|
|
)
|
|
error_code = _safe_error_code(exc, fallback="AGENT_RUNTIME_ERROR")
|
|
error_details = _run_failure_details(exc, error_code)
|
|
if self._last_model_failure is not None:
|
|
error_code, error_details = self._last_model_failure
|
|
await self._terminal_locked(
|
|
"failed",
|
|
error_code=error_code,
|
|
error_details=error_details,
|
|
)
|
|
|
|
async def _stop_checkpoint_writer(self) -> None:
|
|
while True:
|
|
try:
|
|
if self._event_iterator is not None:
|
|
await self._event_iterator.aclose()
|
|
if self._stream_stop is not None:
|
|
await asyncio.shield(self._stream_stop.start([]))
|
|
from ..stream.stop import sqlite_barrier
|
|
await sqlite_barrier(self._host.checkpointer)
|
|
self._checkpoint_writes_stopped = True
|
|
return
|
|
except Exception as exc:
|
|
self._stop_errors.append(exc)
|
|
await asyncio.sleep(0.05)
|
|
|
|
async def _run_title(self) -> None:
|
|
if self._model_set.title is None:
|
|
await self._append_locked(
|
|
"title", {"kind": "skipped", "reason": "TITLE_RUNTIME_INVARIANT"}
|
|
)
|
|
return
|
|
try:
|
|
async with asyncio.timeout(self._snapshot.title_start_timeout_seconds):
|
|
response = await self._model_set.title.ainvoke(
|
|
_title_prompt(str(self._input.message)),
|
|
config={"callbacks": [self._attempt_callback]},
|
|
)
|
|
title = _response_text(response).strip().strip('"').strip()[:100]
|
|
await self._append_locked("title", {"kind": "generated", "title": title})
|
|
except Exception as exc:
|
|
await self._append_locked(
|
|
"title", {"kind": "skipped", "reason": _safe_error_code(exc)}
|
|
)
|
|
|
|
async def _begin_callback_attempt(
|
|
self,
|
|
*,
|
|
callback_run_id: str,
|
|
purpose: str,
|
|
route: ResolvedRoute,
|
|
provider_input_bound_tokens: int,
|
|
provider_input_breakdown: Mapping[str, int] | None = None,
|
|
) -> None:
|
|
async with self._budget_lock:
|
|
attempt_index = self._attempt_counts[purpose] + 1
|
|
bound = next(
|
|
item
|
|
for item in self._snapshot.purpose_route_call_bounds[purpose]
|
|
if item.route_identity == route.identity
|
|
)
|
|
provider_health_state = self._runtime._provider_health.status(
|
|
route.provider_capacity_key
|
|
)
|
|
health_state = self._runtime._route_health.status(route.identity.route_key)
|
|
if (
|
|
provider_health_state == "open"
|
|
or health_state == "open"
|
|
or not self._runtime._route_health.acquire_half_open(
|
|
route.identity.route_key
|
|
)
|
|
):
|
|
self._attempt_counts[purpose] = attempt_index
|
|
rejected = ModelAttemptEvent(
|
|
request_id=self._admission.request_id,
|
|
turn_id=self._admission.turn_id,
|
|
run_id=self._run_id,
|
|
preparation_id=self._admission.preparation_id,
|
|
admission_snapshot_id=self._admission.admission_snapshot_id,
|
|
admission_id=self._admission.admission_id,
|
|
hold_id=self._admission.hold_id,
|
|
prepared_snapshot_digest=self._admission.prepared_snapshot_digest,
|
|
prepared_input_digest=self._admission.prepared_input_digest,
|
|
billing_fencing_token=self._admission.billing_fencing_token,
|
|
turn_fencing_token=self._admission.turn_fencing_token,
|
|
logical_call_id=str(uuid.uuid4()),
|
|
attempt_id=str(uuid.uuid4()),
|
|
purpose=purpose,
|
|
attempt_index=attempt_index,
|
|
identity=route.identity,
|
|
quote_id=route.quote.quote_id,
|
|
billing_intent="user_charge"
|
|
if purpose == "main_agent"
|
|
else "platform_cost",
|
|
outcome="rejected",
|
|
provider_request_started=False,
|
|
provider_input_bound_tokens=provider_input_bound_tokens,
|
|
provider_reserved_microunits=0,
|
|
health_state=health_state,
|
|
error_code="ROUTE_HEALTH_UNAVAILABLE",
|
|
)
|
|
try:
|
|
await self._append_locked("model_attempt", event_payload(rejected))
|
|
except Exception:
|
|
self._attempt_counts[purpose] -= 1
|
|
raise
|
|
raise EvoRuntimeError("ROUTE_HEALTH_UNAVAILABLE")
|
|
if provider_input_bound_tokens > bound.payload_input_hard_cap:
|
|
raise EvoRuntimeError(
|
|
"MODEL_CONTEXT_WINDOW_EXCEEDED",
|
|
details=(
|
|
{
|
|
"provider_input_bound_tokens": provider_input_bound_tokens,
|
|
"payload_input_hard_cap": bound.payload_input_hard_cap,
|
|
**dict(provider_input_breakdown or {}),
|
|
},
|
|
),
|
|
)
|
|
# Usage remains observable and is settled from the terminal event,
|
|
# but there is intentionally no per-run reserve or call-count gate.
|
|
reserved = 0
|
|
self._attempt_counts[purpose] = attempt_index
|
|
logical_call_id = str(uuid.uuid4())
|
|
attempt_id = str(uuid.uuid4())
|
|
attempt = ModelAttemptEvent(
|
|
request_id=self._admission.request_id,
|
|
turn_id=self._admission.turn_id,
|
|
run_id=self._run_id,
|
|
preparation_id=self._admission.preparation_id,
|
|
admission_snapshot_id=self._admission.admission_snapshot_id,
|
|
admission_id=self._admission.admission_id,
|
|
hold_id=self._admission.hold_id,
|
|
prepared_snapshot_digest=self._admission.prepared_snapshot_digest,
|
|
prepared_input_digest=self._admission.prepared_input_digest,
|
|
billing_fencing_token=self._admission.billing_fencing_token,
|
|
turn_fencing_token=self._admission.turn_fencing_token,
|
|
logical_call_id=logical_call_id,
|
|
attempt_id=attempt_id,
|
|
purpose=purpose,
|
|
attempt_index=attempt_index,
|
|
identity=route.identity,
|
|
quote_id=route.quote.quote_id,
|
|
billing_intent="user_charge"
|
|
if purpose == "main_agent"
|
|
else "platform_cost",
|
|
outcome="started",
|
|
provider_request_started=True,
|
|
provider_input_bound_tokens=provider_input_bound_tokens,
|
|
provider_reserved_microunits=reserved,
|
|
health_state=health_state,
|
|
)
|
|
self._callback_attempts[callback_run_id] = (
|
|
attempt,
|
|
reserved,
|
|
route.quote,
|
|
route,
|
|
)
|
|
try:
|
|
await self._append_locked("model_attempt", event_payload(attempt))
|
|
except Exception:
|
|
self._callback_attempts.pop(callback_run_id, None)
|
|
self._remaining += reserved
|
|
self._attempt_counts[purpose] -= 1
|
|
self._runtime._route_health.release_half_open(route.identity.route_key)
|
|
raise
|
|
|
|
async def _finish_callback_attempt(
|
|
self,
|
|
callback_run_id: str,
|
|
*,
|
|
usage: Mapping[str, int | str | None] | None,
|
|
error_code: str | None,
|
|
health_scope: str = "model_route",
|
|
) -> None:
|
|
async with self._budget_lock:
|
|
record = self._callback_attempts.pop(callback_run_id, None)
|
|
if record is None:
|
|
return
|
|
started, _reserved, quote, route = record
|
|
valid_usage = _normalize_usage(usage)
|
|
outcome = "failed" if error_code else "succeeded"
|
|
if valid_usage is not None:
|
|
_actual_cost(valid_usage, quote)
|
|
bound = next(
|
|
item
|
|
for item in self._snapshot.purpose_route_call_bounds[
|
|
started.purpose
|
|
]
|
|
if item.route_identity == started.identity
|
|
)
|
|
if (
|
|
int(valid_usage["input_tokens"]) > bound.billable_input_cap
|
|
or int(valid_usage["output_tokens"]) > bound.max_output_tokens
|
|
):
|
|
valid_usage = None
|
|
error_code = "USAGE_LIMIT_VIOLATION"
|
|
outcome = "failed"
|
|
if valid_usage is None:
|
|
outcome = "usage_unconfirmed"
|
|
terminal = replace(
|
|
started,
|
|
outcome=outcome,
|
|
usage=valid_usage,
|
|
usage_available=valid_usage is not None,
|
|
error_code=error_code,
|
|
timestamp=now_ms(),
|
|
)
|
|
if error_code:
|
|
if health_scope == "provider_connection":
|
|
self._runtime._provider_health.record_failure(
|
|
route.provider_capacity_key, error_code
|
|
)
|
|
elif health_scope == "model_route":
|
|
self._runtime._route_health.record_failure(
|
|
started.identity.route_key, error_code
|
|
)
|
|
else:
|
|
self._runtime._route_health.record_success(started.identity.route_key)
|
|
self._runtime._provider_health.record_success(
|
|
route.provider_capacity_key
|
|
)
|
|
await self._append_locked("model_attempt", event_payload(terminal))
|
|
self._runtime._record_runtime_observation(
|
|
route,
|
|
purpose=started.purpose,
|
|
outcome=outcome,
|
|
error_code=error_code,
|
|
)
|
|
|
|
async def _append_locked(
|
|
self,
|
|
kind: str,
|
|
payload: Mapping[str, Any],
|
|
*,
|
|
control_terminal: bool = False,
|
|
) -> EvoRuntimeEvent:
|
|
async with self._event_commit_lock:
|
|
if self._terminal_event is not None and kind != "run":
|
|
return self._terminal_event
|
|
sequence = self._sequence + 1
|
|
event = EvoRuntimeEvent(
|
|
event_id=str(uuid.uuid4()),
|
|
runtime_instance_id=self._runtime.runtime_instance_id,
|
|
run_id=self._run_id,
|
|
request_id=self._admission.request_id,
|
|
turn_id=self._admission.turn_id,
|
|
sequence=sequence,
|
|
kind=kind,
|
|
payload=dict(payload),
|
|
)
|
|
event_size = len(canonical_json_v1(asdict(event)))
|
|
if not control_terminal and (
|
|
sequence >= self._snapshot.max_run_journal_events
|
|
or self._journal_bytes + event_size
|
|
> self._snapshot.max_run_journal_bytes
|
|
):
|
|
raise EvoRuntimeError("RUN_EVENT_LIMIT_EXCEEDED")
|
|
sink = self._host.runtime_event_sink
|
|
assert sink is not None
|
|
payload_digest = hashlib.sha256(
|
|
canonical_json_v1(asdict(event))
|
|
).hexdigest()
|
|
registry = self._runtime.host_registry
|
|
if control_terminal and registry is not None:
|
|
await asyncio.to_thread(registry.prepare_terminal, self.run_id, event=asdict(event))
|
|
commit_error: Exception | None = None
|
|
try:
|
|
result = await sink.commit(event)
|
|
except EvoRuntimeError:
|
|
raise
|
|
except Exception as exc:
|
|
commit_error = exc
|
|
result = "unknown"
|
|
if result not in {"committed", "duplicate"}:
|
|
try:
|
|
confirmation = await sink.confirm(event.event_id, payload_digest)
|
|
except Exception as exc:
|
|
raise EvoRuntimeError("EVENT_COMMIT_INDETERMINATE") from exc
|
|
if confirmation == "committed":
|
|
result = "committed"
|
|
elif confirmation == "absent":
|
|
raise EvoRuntimeError("EVENT_INGRESS_UNAVAILABLE") from commit_error
|
|
elif confirmation == "conflict":
|
|
raise EvoRuntimeError("EVO_EVENT_CONFLICT")
|
|
else:
|
|
raise EvoRuntimeError("EVENT_COMMIT_INDETERMINATE")
|
|
if control_terminal and registry is not None:
|
|
await asyncio.to_thread(registry.confirm_terminal, self.run_id, digest=payload_digest)
|
|
async with self._condition:
|
|
self._sequence = sequence
|
|
self._journal_bytes += event_size
|
|
self._journal.append(event)
|
|
self._condition.notify_all()
|
|
return event
|
|
|
|
async def _terminal_locked(
|
|
self,
|
|
outcome: str,
|
|
**details: Any,
|
|
) -> EvoRuntimeEvent:
|
|
# Freeze the first terminal intent before yielding. Agent timeout/cancel
|
|
# cannot cancel its compensation or substitute a different outcome.
|
|
return await asyncio.shield(self._start_terminal(outcome, **details))
|
|
|
|
def _start_terminal(self, outcome: str, **details: Any) -> asyncio.Task[EvoRuntimeEvent]:
|
|
if self._terminal_task is None:
|
|
self._terminal_task = asyncio.create_task(
|
|
self._finalize_terminal(outcome, **details),
|
|
name=f"evo-terminal-{self._run_id}",
|
|
)
|
|
self._terminal_task.add_done_callback(self._observe_terminal)
|
|
return self._terminal_task
|
|
|
|
def _observe_terminal(self, task: asyncio.Task[EvoRuntimeEvent]) -> None:
|
|
error = None if task.cancelled() else task.exception()
|
|
if task.cancelled() or error is not None:
|
|
self._terminal_error_code = (
|
|
"RUN_TERMINAL_CANCELLED" if task.cancelled() else
|
|
"RUN_TERMINAL_CLEANUP_UNCONFIRMED" if self._terminal_phase == "cleanup" else
|
|
"RUN_TERMINAL_COMMIT_UNCONFIRMED")
|
|
logger.error("owner terminal task failed: %s", self._terminal_error_code)
|
|
|
|
async def _finish_registry(self, outcome: str, checkpoint_id: str = "") -> None:
|
|
delay = 0.1
|
|
while self._runtime.host_registry is not None:
|
|
try:
|
|
await asyncio.to_thread(
|
|
self._runtime.host_registry.finish, self.run_id,
|
|
outcome=outcome, checkpoint_id=checkpoint_id,
|
|
)
|
|
return
|
|
except Exception:
|
|
# Product-owned compensation, never an Agent retry. Keep the
|
|
# claim and fixed outcome until the original owner confirms it.
|
|
logger.warning("registry finalization pending for %s", self.run_id)
|
|
await asyncio.sleep(delay)
|
|
delay = min(delay * 2, 5.0)
|
|
|
|
async def _finalize_terminal(
|
|
self,
|
|
outcome: str,
|
|
*,
|
|
error_code: str | None = None,
|
|
reason: str | None = None,
|
|
error_details: Mapping[str, Any] | None = None,
|
|
checkpoint_details: Mapping[str, Any] | None = None,
|
|
) -> EvoRuntimeEvent:
|
|
# Terminalization is run-owned. Only public observers have deadlines;
|
|
# a slow close must not permanently fail the finalization task.
|
|
if self._stop_task is not None:
|
|
await asyncio.shield(self._stop_task)
|
|
await self._ensure_cleanup(timeout=None)
|
|
if self._runtime.host_registry is not None and not self._checkpoint_writes_stopped:
|
|
raise EvoRuntimeError("CHECKPOINT_WRITER_STOP_UNCONFIRMED")
|
|
self._terminal_phase = "commit"
|
|
async with self._terminal_lock:
|
|
if self._terminal_event is not None:
|
|
return self._terminal_event
|
|
event = await self._append_locked(
|
|
"run",
|
|
{
|
|
"kind": "run_terminal",
|
|
"preparation_id": self._admission.preparation_id,
|
|
"admission_snapshot_id": self._admission.admission_snapshot_id,
|
|
"admission_id": self._admission.admission_id,
|
|
"hold_id": self._admission.hold_id,
|
|
"turn_fencing_token": self._admission.turn_fencing_token,
|
|
"billing_fencing_token": self._admission.billing_fencing_token,
|
|
"outcome": outcome,
|
|
"error_code": error_code,
|
|
"reason": reason,
|
|
"error_details": dict(error_details or {}),
|
|
"purpose_attempt_counts": dict(self._attempt_counts),
|
|
"provider_reserve_remaining_microunits": self._remaining,
|
|
**dict(checkpoint_details or {}),
|
|
},
|
|
control_terminal=True,
|
|
)
|
|
await self._finish_registry(
|
|
outcome, str((checkpoint_details or {}).get("checkpoint_id") or ""))
|
|
if self._prepared.state == "CONSTRUCTION_FAILED":
|
|
async with self._runtime._registry_lock:
|
|
self._runtime._retire_prepared_locked(self._prepared)
|
|
async with self._condition:
|
|
self._state = "TERMINAL"
|
|
self._terminal_event = event
|
|
self._condition.notify_all()
|
|
return event
|
|
|
|
async def _replay(self, cursor: int) -> AsyncIterator[EvoRuntimeEvent]:
|
|
while True:
|
|
async with self._condition:
|
|
while self._sequence <= cursor and self._terminal_event is None:
|
|
if self._agent_task is not None and self._agent_task.done():
|
|
break
|
|
await self._condition.wait()
|
|
events = [event for event in self._journal if event.sequence > cursor]
|
|
terminal = self._terminal_event is not None
|
|
agent_task = self._agent_task
|
|
for event in events:
|
|
cursor = event.sequence
|
|
yield event
|
|
if terminal and cursor >= self._sequence:
|
|
return
|
|
if not events and agent_task is not None and agent_task.done():
|
|
error = agent_task.exception()
|
|
if error is not None:
|
|
raise error
|
|
raise EvoRuntimeError("EVENT_REPLAY_UNAVAILABLE")
|
|
|
|
|
|
class _RuntimeAttemptCallback(AsyncCallbackHandler):
|
|
# LangChain otherwise logs and suppresses callback failures, which would let
|
|
# a Provider request start without a committed STARTED attempt event.
|
|
raise_error = True
|
|
|
|
def __init__(self, run: _EvoWebRun) -> None:
|
|
self.run = run
|
|
self._stream_diagnostics: dict[str, _ModelStreamDiagnosticState] = {}
|
|
|
|
async def on_chat_model_start(
|
|
self,
|
|
_serialized: dict[str, Any],
|
|
messages: list[list[Any]],
|
|
*,
|
|
run_id: Any,
|
|
metadata: Mapping[str, Any] | None = None,
|
|
**_kwargs: Any,
|
|
) -> None:
|
|
purpose, route = self._route_for(metadata)
|
|
try:
|
|
callback_payload = _callback_messages_payload(messages)
|
|
# Structured output and bound tools live in invocation params, not messages.
|
|
invocation = _kwargs.get("invocation_params", {})
|
|
input_parameters = _effective_callback_input_parameters(invocation)
|
|
input_bound = _provider_input_token_bound({
|
|
"messages": callback_payload,
|
|
**input_parameters,
|
|
})
|
|
if purpose == "main_agent" and self.run._force_context_repair_pending:
|
|
self.run._force_context_repair_pending = False
|
|
raise EvoRuntimeError(
|
|
"MODEL_CONTEXT_WINDOW_EXCEEDED",
|
|
details=(
|
|
{
|
|
"reason": "forced_context_repair",
|
|
"repair_requested": True,
|
|
**input_bound.projection(),
|
|
},
|
|
),
|
|
)
|
|
await self.run._begin_callback_attempt(
|
|
callback_run_id=str(run_id),
|
|
purpose=purpose,
|
|
route=route,
|
|
provider_input_bound_tokens=input_bound.total_tokens,
|
|
provider_input_breakdown=input_bound.projection(),
|
|
)
|
|
summary = _callback_payload_debug_summary(
|
|
callback_payload, input_bound.total_tokens
|
|
)
|
|
plan = route.invocation_plan
|
|
plan_params = ",".join(
|
|
sorted(str(key)[:64] for key in (plan.sdk_params if plan else {}))
|
|
)
|
|
parameter_values = _invocation_parameters_debug(plan)
|
|
message_schema = _callback_message_schema_debug(callback_payload)
|
|
attempt_record = self.run._callback_attempts.get(str(run_id))
|
|
attempt = attempt_record[0] if attempt_record is not None else None
|
|
started_at = time.monotonic()
|
|
self._stream_diagnostics[str(run_id)] = _ModelStreamDiagnosticState(
|
|
started_at=started_at,
|
|
last_checkpoint_at=started_at,
|
|
attempt_id=attempt.attempt_id if attempt is not None else "",
|
|
purpose=purpose,
|
|
route_key=route.identity.route_key,
|
|
)
|
|
logger.info(
|
|
"model_request_debug call_id=%s attempt_id=%s attempt_index=%s "
|
|
"purpose=%s route=%s plan_hash=%s input_bound_tokens=%s "
|
|
"input_bytes=%s input_sha256=%s "
|
|
"api_mode=%s tool_transport=%s output=%s=%s reasoning=%s "
|
|
"sdk_param_keys=%s params=%s message_types=%s content_block_types=%s "
|
|
"tool_calls=%s tool_results=%s message_schema=%s",
|
|
str(run_id),
|
|
attempt.attempt_id if attempt is not None else "",
|
|
attempt.attempt_index if attempt is not None else "",
|
|
purpose,
|
|
route.identity.route_key,
|
|
plan.plan_hash if plan else "",
|
|
summary["input_bound_tokens"],
|
|
summary["input_bytes"],
|
|
summary["input_sha256"],
|
|
plan.api_mode if plan else route.identity.api_mode,
|
|
plan.tool_call_transport
|
|
if plan
|
|
else route.identity.tool_call_transport,
|
|
plan.output_token_parameter if plan else "",
|
|
plan.output_token_limit if plan else route.max_output_tokens,
|
|
plan.reasoning_effort if plan else "",
|
|
plan_params,
|
|
parameter_values,
|
|
summary["message_types"],
|
|
summary["content_block_types"],
|
|
summary["tool_calls"],
|
|
summary["tool_results"],
|
|
message_schema,
|
|
)
|
|
logger.info(
|
|
"model_stream_checkpoint phase=request_started call_id=%s "
|
|
"attempt_id=%s purpose=%s route=%s plan_hash=%s "
|
|
"prepared_input_digest=%s visible_chunks=0",
|
|
str(run_id),
|
|
attempt.attempt_id if attempt is not None else "",
|
|
purpose,
|
|
route.identity.route_key,
|
|
plan.plan_hash if plan else "",
|
|
attempt.prepared_input_digest if attempt is not None else "",
|
|
)
|
|
except EvoRuntimeError as exc:
|
|
self.run._last_model_failure = (
|
|
exc.code,
|
|
_callback_start_failure_details(purpose, route, exc),
|
|
)
|
|
raise
|
|
|
|
async def on_llm_new_token(
|
|
self,
|
|
token: str | list[str | dict[str, Any]],
|
|
*,
|
|
chunk: Any = None,
|
|
run_id: Any,
|
|
**_kwargs: Any,
|
|
) -> None:
|
|
state = self._stream_diagnostics.get(str(run_id))
|
|
if state is None:
|
|
return
|
|
now = time.monotonic()
|
|
previous_visible_at = state.last_visible_at
|
|
state.visible_chunks += 1
|
|
state.text_chars += _stream_token_chars(token)
|
|
kinds = _stream_chunk_kinds(token, chunk)
|
|
for kind in kinds:
|
|
state.kind_counts[kind] = state.kind_counts.get(kind, 0) + 1
|
|
state.last_visible_at = now
|
|
should_log = (
|
|
state.visible_chunks in {1, 10, 100}
|
|
or state.visible_chunks % 1_000 == 0
|
|
or now - state.last_checkpoint_at >= 30.0
|
|
)
|
|
if not should_log:
|
|
return
|
|
phase = (
|
|
"first_visible_chunk" if state.visible_chunks == 1 else "stream_progress"
|
|
)
|
|
idle_before_ms = int(
|
|
max(0.0, now - (previous_visible_at or state.started_at)) * 1_000
|
|
)
|
|
state.last_checkpoint_at = now
|
|
logger.info(
|
|
"model_stream_checkpoint phase=%s call_id=%s attempt_id=%s purpose=%s "
|
|
"route=%s elapsed_ms=%s idle_before_ms=%s visible_chunks=%s "
|
|
"text_chars=%s chunk_kinds=%s",
|
|
phase,
|
|
str(run_id),
|
|
state.attempt_id,
|
|
state.purpose,
|
|
state.route_key,
|
|
int(max(0.0, now - state.started_at) * 1_000),
|
|
idle_before_ms,
|
|
state.visible_chunks,
|
|
state.text_chars,
|
|
_stream_kind_counts_debug(state.kind_counts),
|
|
)
|
|
|
|
async def on_llm_end(self, response: Any, *, run_id: Any, **_kwargs: Any) -> None:
|
|
self._log_terminal_checkpoint(str(run_id), phase="completed")
|
|
self.run._last_model_failure = None
|
|
await self.run._finish_callback_attempt(
|
|
str(run_id), usage=_usage_from_response(response), error_code=None
|
|
)
|
|
|
|
async def on_llm_error(
|
|
self, error: BaseException, *, run_id: Any, **_kwargs: Any
|
|
) -> None:
|
|
record = self.run._callback_attempts.get(str(run_id))
|
|
route = record[3] if record is not None else None
|
|
disposition = (
|
|
route.adapter.classify_error(error)
|
|
if route is not None and route.adapter is not None
|
|
else None
|
|
)
|
|
error_code = disposition.error_code if disposition else _safe_error_code(error)
|
|
parse_debug = _provider_json_decode_debug(error)
|
|
self._log_terminal_checkpoint(str(run_id), phase="error", error=error)
|
|
if parse_debug:
|
|
logger.warning(
|
|
"model_response_json_decode purpose=%s route=%s document_bytes=%s "
|
|
"document_sha256=%s line=%s column=%s position=%s "
|
|
"invalid_codepoint=%s inside_string=%s event_type=%s "
|
|
"window_bytes=%s window_sha256=%s control_codepoints=%s trace=%s",
|
|
record[0].purpose if record is not None else "",
|
|
route.identity.route_key if route is not None else "",
|
|
parse_debug["provider_json_document_bytes"],
|
|
parse_debug["provider_json_document_sha256"],
|
|
parse_debug["provider_json_line"],
|
|
parse_debug["provider_json_column"],
|
|
parse_debug["provider_json_position"],
|
|
parse_debug.get("provider_json_invalid_codepoint", ""),
|
|
parse_debug["provider_json_position_inside_string"],
|
|
parse_debug.get("provider_json_event_type", ""),
|
|
parse_debug["provider_json_window_bytes"],
|
|
parse_debug["provider_json_window_sha256"],
|
|
parse_debug.get("provider_json_control_codepoints", ""),
|
|
parse_debug["provider_json_trace"],
|
|
)
|
|
failure_details = _provider_failure_details(error, route, error_code)
|
|
provider_message = _provider_error_message_debug(error, route)
|
|
logger.warning(
|
|
"model_provider_failure purpose=%s route=%s code=%s status=%s "
|
|
"provider_code=%s provider_param=%s api_mode=%s tool_transport=%s "
|
|
"output=%s=%s provider_message=%s",
|
|
record[0].purpose if record is not None else "",
|
|
route.identity.route_key if route is not None else "",
|
|
error_code,
|
|
failure_details.get("http_status", ""),
|
|
failure_details.get("provider_error_code", ""),
|
|
failure_details.get("provider_error_parameter", ""),
|
|
failure_details.get("api_mode", ""),
|
|
failure_details.get("tool_call_transport", ""),
|
|
failure_details.get("output_token_parameter", ""),
|
|
failure_details.get("output_token_limit", ""),
|
|
provider_message,
|
|
)
|
|
self.run._last_model_failure = (error_code, failure_details)
|
|
await self.run._finish_callback_attempt(
|
|
str(run_id),
|
|
usage=None,
|
|
error_code=error_code,
|
|
health_scope=(disposition.health_scope if disposition else "model_route"),
|
|
)
|
|
|
|
def _log_terminal_checkpoint(
|
|
self, run_id: str, *, phase: str, error: BaseException | None = None
|
|
) -> None:
|
|
state = self._stream_diagnostics.pop(run_id, None)
|
|
if state is None:
|
|
return
|
|
now = time.monotonic()
|
|
timeout_debug = _provider_stream_timeout_debug(error)
|
|
last_activity = state.last_visible_at or state.started_at
|
|
log = logger.warning if error is not None else logger.info
|
|
log(
|
|
"model_stream_checkpoint phase=%s call_id=%s attempt_id=%s purpose=%s "
|
|
"route=%s elapsed_ms=%s idle_since_visible_ms=%s visible_chunks=%s "
|
|
"text_chars=%s chunk_kinds=%s provider_chunks_received=%s "
|
|
"provider_idle_timeout_ms=%s error_type=%s",
|
|
phase,
|
|
run_id,
|
|
state.attempt_id,
|
|
state.purpose,
|
|
state.route_key,
|
|
int(max(0.0, now - state.started_at) * 1_000),
|
|
int(max(0.0, now - last_activity) * 1_000),
|
|
state.visible_chunks,
|
|
state.text_chars,
|
|
_stream_kind_counts_debug(state.kind_counts),
|
|
timeout_debug.get("provider_stream_chunks_received", ""),
|
|
timeout_debug.get("provider_stream_idle_timeout_ms", ""),
|
|
type(error).__name__[:128] if error is not None else "",
|
|
)
|
|
|
|
def _route_for(
|
|
self, metadata: Mapping[str, Any] | None
|
|
) -> tuple[str, ResolvedRoute]:
|
|
values = dict(metadata or {})
|
|
purpose = str(values.get("runtime_purpose") or "main_agent")
|
|
route_key = str(values.get("route_key") or "")
|
|
routes = self.run._snapshot.purpose_routes.get(purpose, ())
|
|
for route in routes:
|
|
if route.identity.route_key == route_key or (
|
|
not route_key and len(routes) == 1
|
|
):
|
|
return purpose, route
|
|
raise EvoRuntimeError("ADMISSION_STALE")
|
|
|
|
|
|
@dataclass(slots=True)
|
|
class _ModelStreamDiagnosticState:
|
|
started_at: float
|
|
last_checkpoint_at: float
|
|
attempt_id: str
|
|
purpose: str
|
|
route_key: str
|
|
visible_chunks: int = 0
|
|
text_chars: int = 0
|
|
last_visible_at: float | None = None
|
|
kind_counts: dict[str, int] = field(default_factory=dict)
|
|
|
|
|
|
def _stream_token_chars(token: Any) -> int:
|
|
if isinstance(token, str):
|
|
return len(token)
|
|
if not isinstance(token, list):
|
|
return 0
|
|
return sum(
|
|
len(str(item.get("text") or "")) for item in token if isinstance(item, Mapping)
|
|
)
|
|
|
|
|
|
def _stream_chunk_kinds(token: Any, chunk: Any) -> tuple[str, ...]:
|
|
kinds: set[str] = set()
|
|
if _stream_token_chars(token) > 0:
|
|
kinds.add("text")
|
|
message = getattr(chunk, "message", None)
|
|
content = getattr(message, "content", None)
|
|
if isinstance(content, list):
|
|
for block in content:
|
|
if not isinstance(block, Mapping):
|
|
continue
|
|
block_type = str(block.get("type") or "")
|
|
if "reasoning" in block_type:
|
|
kinds.add("reasoning")
|
|
elif block_type in {"text", "output_text"}:
|
|
kinds.add("text")
|
|
elif "tool" in block_type or "function" in block_type:
|
|
kinds.add("tool_call")
|
|
additional = getattr(message, "additional_kwargs", None)
|
|
if isinstance(additional, Mapping) and any(
|
|
key in additional for key in ("reasoning", "reasoning_content")
|
|
):
|
|
kinds.add("reasoning")
|
|
if getattr(message, "tool_call_chunks", None):
|
|
kinds.add("tool_call")
|
|
return tuple(sorted(kinds or {"metadata"}))
|
|
|
|
|
|
def _stream_kind_counts_debug(counts: Mapping[str, int]) -> str:
|
|
return ",".join(
|
|
f"{kind}:{int(counts.get(kind, 0))}"
|
|
for kind in ("text", "reasoning", "tool_call", "metadata")
|
|
if counts.get(kind, 0)
|
|
)
|
|
|
|
|
|
def _callback_start_failure_details(
|
|
purpose: str, route: ResolvedRoute, error: EvoRuntimeError
|
|
) -> dict[str, str | int]:
|
|
"""Expose a bounded callback-start failure without provider contents."""
|
|
|
|
plan = route.invocation_plan
|
|
details: dict[str, str | int] = {
|
|
"failure_stage": "model_attempt_start",
|
|
"reason": error.code.lower(),
|
|
"purpose": purpose,
|
|
"provider": route.identity.provider_id,
|
|
"model": route.identity.model_id,
|
|
"route_key": route.identity.route_key,
|
|
"endpoint": route.identity.endpoint_name,
|
|
"api_mode": plan.api_mode if plan else route.identity.api_mode,
|
|
"tool_call_transport": (
|
|
plan.tool_call_transport if plan else route.identity.tool_call_transport
|
|
),
|
|
}
|
|
if plan is not None:
|
|
details["output_token_limit"] = plan.output_token_limit
|
|
details["output_token_parameter"] = plan.output_token_parameter
|
|
details["invocation_plan_hash"] = plan.plan_hash
|
|
for item in error.details:
|
|
for key, value in item.items():
|
|
if isinstance(value, str | int | bool):
|
|
details[str(key)] = value
|
|
return details
|
|
|
|
|
|
def _checked_ceil_cost(
|
|
input_tokens: int,
|
|
input_rate: int,
|
|
output_tokens: int,
|
|
output_rate: int,
|
|
unit_scale: int,
|
|
) -> int:
|
|
numerator = input_tokens * input_rate + output_tokens * output_rate
|
|
if numerator < 0 or numerator > _BIGINT_MAX * unit_scale:
|
|
raise EvoRuntimeError("COST_OVERFLOW")
|
|
result = math.ceil(numerator / unit_scale)
|
|
if result > _BIGINT_MAX:
|
|
raise EvoRuntimeError("COST_OVERFLOW")
|
|
return result
|
|
|
|
|
|
def _payload_token_bound(payload: Any) -> int:
|
|
try:
|
|
return len(canonical_json_v1(payload))
|
|
except Exception as exc:
|
|
raise EvoRuntimeError("TOKEN_BOUND_UNAVAILABLE") from exc
|
|
|
|
|
|
@dataclass(frozen=True, slots=True)
|
|
class _ProviderInputBound:
|
|
total_tokens: int
|
|
text_tokens: int
|
|
media_tokens: int
|
|
media_blocks: int
|
|
largest_media_bytes: int
|
|
|
|
def projection(self) -> dict[str, int]:
|
|
return {
|
|
"provider_input_bound_tokens": self.total_tokens,
|
|
"text_input_bound_tokens": self.text_tokens,
|
|
"media_input_bound_tokens": self.media_tokens,
|
|
"media_blocks": self.media_blocks,
|
|
"largest_media_bytes": self.largest_media_bytes,
|
|
}
|
|
|
|
|
|
def _decode_media_payload(payload: str, mime: str) -> tuple[bytes, str]:
|
|
if payload.startswith("data:"):
|
|
# Bound the encoded allocation before split/decode using the same image
|
|
# ingress limit, not the ordinary JSON/schema text budget.
|
|
from EvoScientist.document_extract import MAX_IMAGE_BYTES
|
|
|
|
if len(payload) > 4 * ((MAX_IMAGE_BYTES + 2) // 3) + 256:
|
|
raise EvoRuntimeError("TOKEN_BOUND_UNAVAILABLE")
|
|
try:
|
|
header, encoded = payload.split(",", 1)
|
|
resolved_mime = header[5:].split(";", 1)[0] or mime
|
|
if ";base64" not in header.lower():
|
|
raise ValueError("media data URL is not base64 encoded")
|
|
return base64.b64decode(encoded, validate=True), resolved_mime
|
|
except (ValueError, binascii.Error) as exc:
|
|
raise EvoRuntimeError("TOKEN_BOUND_UNAVAILABLE") from exc
|
|
try:
|
|
return base64.b64decode(payload, validate=True), mime
|
|
except (ValueError, binascii.Error) as exc:
|
|
raise EvoRuntimeError("TOKEN_BOUND_UNAVAILABLE") from exc
|
|
|
|
|
|
def _bounded_media_block(
|
|
value: Mapping[str, Any],
|
|
) -> tuple[dict[str, Any], int, int] | None:
|
|
block = dict(value)
|
|
mime = str(block.get("mime_type") or "application/octet-stream")
|
|
payload: str | None = None
|
|
location: tuple[str, ...] = ()
|
|
|
|
if isinstance(block.get("base64"), str):
|
|
payload = str(block["base64"])
|
|
location = ("base64",)
|
|
elif isinstance(block.get("url"), str) and str(block["url"]).startswith("data:"):
|
|
payload = str(block["url"])
|
|
location = ("url",)
|
|
elif isinstance(block.get("image_url"), str) and str(block["image_url"]).startswith(
|
|
"data:"
|
|
):
|
|
payload = str(block["image_url"])
|
|
location = ("image_url",)
|
|
elif isinstance(block.get("image_url"), Mapping):
|
|
image_url = dict(block["image_url"])
|
|
if isinstance(image_url.get("url"), str) and str(image_url["url"]).startswith(
|
|
"data:"
|
|
):
|
|
payload = str(image_url["url"])
|
|
location = ("image_url", "url")
|
|
elif isinstance(block.get("source"), Mapping):
|
|
source = dict(block["source"])
|
|
if source.get("type") == "base64" and isinstance(source.get("data"), str):
|
|
payload = str(source["data"])
|
|
mime = str(source.get("media_type") or mime)
|
|
location = ("source", "data")
|
|
elif isinstance(block.get("inline_data"), Mapping):
|
|
inline_data = dict(block["inline_data"])
|
|
if isinstance(inline_data.get("data"), str):
|
|
payload = str(inline_data["data"])
|
|
mime = str(inline_data.get("mime_type") or mime)
|
|
location = ("inline_data", "data")
|
|
|
|
if payload is None:
|
|
return None
|
|
raw, mime = _decode_media_payload(payload, mime)
|
|
digest = hashlib.sha256(raw).hexdigest()[:24]
|
|
marker = f"<media:{mime}:sha256:{digest}>"
|
|
if location == ("base64",):
|
|
block["base64"] = marker
|
|
elif location == ("url",):
|
|
block["url"] = marker
|
|
elif location == ("image_url",):
|
|
block["image_url"] = marker
|
|
elif location == ("image_url", "url"):
|
|
nested = dict(block["image_url"])
|
|
nested["url"] = marker
|
|
block["image_url"] = nested
|
|
elif location == ("source", "data"):
|
|
nested = dict(block["source"])
|
|
nested["data"] = marker
|
|
block["source"] = nested
|
|
elif location == ("inline_data", "data"):
|
|
nested = dict(block["inline_data"])
|
|
nested["data"] = marker
|
|
block["inline_data"] = nested
|
|
# This is the existing conservative media bound, now applied to decoded
|
|
# media bytes rather than to the larger base64/JSON representation.
|
|
return block, (len(raw) + 2) // 3 + 512, len(raw)
|
|
|
|
|
|
def _project_media_for_bound(value: Any, metrics: dict[str, int]) -> Any:
|
|
if isinstance(value, Mapping):
|
|
bounded = _bounded_media_block(value)
|
|
if bounded is not None:
|
|
projected, media_bound, media_bytes = bounded
|
|
metrics["media_tokens"] += media_bound
|
|
metrics["media_blocks"] += 1
|
|
metrics["largest_media_bytes"] = max(
|
|
metrics["largest_media_bytes"], media_bytes
|
|
)
|
|
return {
|
|
key: _project_media_for_bound(item, metrics)
|
|
for key, item in projected.items()
|
|
}
|
|
return {
|
|
key: _project_media_for_bound(item, metrics) for key, item in value.items()
|
|
}
|
|
if isinstance(value, list | tuple):
|
|
return [_project_media_for_bound(item, metrics) for item in value]
|
|
return value
|
|
|
|
|
|
def _provider_input_token_bound(payload: Any) -> _ProviderInputBound:
|
|
metrics = {"media_tokens": 0, "media_blocks": 0, "largest_media_bytes": 0}
|
|
projected = _callback_input_parameters(payload, media_metrics=metrics)
|
|
text_tokens = _payload_token_bound(projected)
|
|
media_tokens = metrics["media_tokens"]
|
|
return _ProviderInputBound(
|
|
total_tokens=text_tokens + media_tokens,
|
|
text_tokens=text_tokens,
|
|
media_tokens=media_tokens,
|
|
media_blocks=metrics["media_blocks"],
|
|
largest_media_bytes=metrics["largest_media_bytes"],
|
|
)
|
|
|
|
|
|
_INPUT_PROJECTION_MAX_DEPTH = 64
|
|
_INPUT_PROJECTION_MAX_NODES = 100_000
|
|
_INPUT_PROJECTION_MAX_BYTES = 8 * 1024 * 1024
|
|
|
|
|
|
def _effective_callback_input_parameters(invocation: Any) -> dict[str, Any]:
|
|
"""Project already-flattened OpenAI callback params, then SDK shallow merge.
|
|
|
|
ChatOpenAI._default_params expands model_kwargs; _get_invocation_params
|
|
overlays call kwargs. OpenAI BaseClient._build_request uses _merge_mappings
|
|
with extra_json last (not a recursive merge). A remaining model_kwargs is
|
|
therefore NOT a second defaults layer we can safely interpret.
|
|
Unknown body extensions and explicit alternative context transports fail closed.
|
|
Responses use_previous_response_id may generate an ID after this callback;
|
|
that SDK-managed context is outside this projection's visibility, not banned.
|
|
"""
|
|
fields = {
|
|
"tools", "functions", "response_format", "tool_choice", "function_call",
|
|
"system", "instructions",
|
|
}
|
|
if type(invocation) is not dict or any(type(key) is not str for key in invocation):
|
|
raise EvoRuntimeError("MODEL_INPUT_PROJECTION_INVALID")
|
|
if any(key in invocation for key in (
|
|
"model_kwargs", "messages", "input", "text", "previous_response_id",
|
|
"conversation", "prompt",
|
|
)):
|
|
raise EvoRuntimeError("MODEL_INPUT_PROJECTION_INVALID")
|
|
extra = invocation.get("extra_body")
|
|
if extra is None:
|
|
extra = {}
|
|
# These non-input controls are emitted by adapter_registry/models. Validate
|
|
# their data without dropping or altering them in the actual SDK request.
|
|
controls = {"enable_thinking", "thinking_budget", "thinking"}
|
|
if type(extra) is not dict or any(
|
|
type(key) is not str or key not in fields | controls for key in extra
|
|
):
|
|
raise EvoRuntimeError("MODEL_INPUT_PROJECTION_INVALID")
|
|
effective = {key: invocation[key] for key in fields if key in invocation}
|
|
effective.update(extra)
|
|
validated = _callback_input_parameters(effective)
|
|
return {key: value for key, value in validated.items() if key in fields}
|
|
|
|
|
|
def _callback_input_parameters(
|
|
value: Any, *, media_metrics: dict[str, int] | None = None
|
|
) -> Any:
|
|
"""Validate JSON/schema trees before serialization; never format rejected values.
|
|
|
|
Limits cover depth (root=0), visited nodes including keys, and aggregate
|
|
JSON scalar/container bytes. Shared acyclic values are counted each time.
|
|
BaseModel classes are trusted application code: schema hooks execute before
|
|
their output can be bounded. This is not a sandbox for untrusted classes.
|
|
"""
|
|
from pydantic import BaseModel
|
|
|
|
active: set[int] = set()
|
|
nodes = 0
|
|
size = 0
|
|
|
|
def visit(item: Any, depth: int, path: tuple[str, ...] = (),
|
|
media: tuple[tuple[str, ...], str] | None = None) -> Any:
|
|
nonlocal nodes, size
|
|
nodes += 1
|
|
if depth > _INPUT_PROJECTION_MAX_DEPTH or nodes > _INPUT_PROJECTION_MAX_NODES:
|
|
raise ValueError
|
|
kind = type(item)
|
|
if item is None or kind in (bool, int, float, str):
|
|
if kind is str and media is not None and path == media[0]:
|
|
from EvoScientist.document_extract import MAX_IMAGE_BYTES, prepare_image_bytes
|
|
|
|
if len(item) > 4 * ((MAX_IMAGE_BYTES + 2) // 3) + 256:
|
|
raise ValueError
|
|
raw, mime = _decode_media_payload(item, media[1])
|
|
if not mime.startswith("image/"):
|
|
raise ValueError
|
|
# Parse using the ingress pixel/byte limits. Do not replace the
|
|
# actual request with the downsampled result or with the marker.
|
|
prepare_image_bytes(raw, "callback-image")
|
|
assert media_metrics is not None
|
|
media_metrics["media_tokens"] += (len(raw) + 2) // 3 + 512
|
|
media_metrics["media_blocks"] += 1
|
|
media_metrics["largest_media_bytes"] = max(
|
|
media_metrics["largest_media_bytes"], len(raw)
|
|
)
|
|
item = f"<media:{mime}:sha256:{hashlib.sha256(raw).hexdigest()[:24]}>"
|
|
elif kind is str and media_metrics is not None and item.startswith("data:"):
|
|
raise ValueError
|
|
if isinstance(item, float) and not math.isfinite(item):
|
|
raise ValueError
|
|
if isinstance(item, str) and len(item) > _INPUT_PROJECTION_MAX_BYTES:
|
|
raise ValueError
|
|
size += len(json.dumps(item, ensure_ascii=True, allow_nan=False))
|
|
if size > _INPUT_PROJECTION_MAX_BYTES:
|
|
raise ValueError
|
|
return item
|
|
is_schema = isinstance(item, type) and issubclass(item, BaseModel)
|
|
if kind not in (dict, list) and not is_schema:
|
|
raise ValueError
|
|
identity = id(item)
|
|
if identity in active:
|
|
raise ValueError
|
|
active.add(identity)
|
|
try:
|
|
if isinstance(item, type) and issubclass(item, BaseModel):
|
|
return visit(item.model_json_schema(), depth + 1)
|
|
if not isinstance(item, (dict, list)):
|
|
raise ValueError
|
|
if len(item) > _INPUT_PROJECTION_MAX_NODES - nodes:
|
|
raise ValueError
|
|
size += 2 + 2 * len(item)
|
|
if size > _INPUT_PROJECTION_MAX_BYTES:
|
|
raise ValueError
|
|
if isinstance(item, list):
|
|
return [visit(child, depth + 1, path + ("[]",)) for child in item]
|
|
# Only content-list image blocks (or a direct block-list caller) can
|
|
# exempt payload bytes. Tools/schema objects never acquire this role.
|
|
content_block = path in {
|
|
("[]",),
|
|
("messages", "[]", "[]", "data", "content", "[]"),
|
|
("messages", "[]", "content", "[]"),
|
|
}
|
|
if media_metrics is not None and content_block:
|
|
block_type = item.get("type")
|
|
mime = item.get("mime_type", "image/unknown")
|
|
if type(item.get("inline_data")) is dict:
|
|
inline_mime = item["inline_data"].get("mime_type")
|
|
if type(inline_mime) is str and inline_mime.startswith("image/"):
|
|
media = (path + ("inline_data", "data"), inline_mime)
|
|
if type(block_type) is str and block_type in {"image", "image_url", "input_image"}:
|
|
location = None
|
|
if type(item.get("base64")) is str:
|
|
location = ("base64",)
|
|
elif (type(item.get("image_url")) is dict
|
|
and type(item["image_url"].get("url")) is str
|
|
and item["image_url"]["url"].startswith("data:")):
|
|
location = ("image_url", "url")
|
|
elif (type(item.get("image_url")) is str
|
|
and item["image_url"].startswith("data:")):
|
|
location = ("image_url",)
|
|
elif type(item.get("url")) is str and item["url"].startswith("data:"):
|
|
location = ("url",)
|
|
elif type(item.get("source")) is dict and item["source"].get("type") == "base64":
|
|
location = ("source", "data")
|
|
mime = item["source"].get("media_type", mime)
|
|
if location is not None:
|
|
if type(mime) is not str:
|
|
raise ValueError
|
|
media = (path + location, mime)
|
|
result = {}
|
|
for key, child in item.items():
|
|
if type(key) is not str:
|
|
raise ValueError
|
|
visit(key, depth + 1)
|
|
result[key] = visit(child, depth + 1, path + (key,), media)
|
|
return result
|
|
finally:
|
|
active.remove(identity)
|
|
|
|
try:
|
|
return visit(value, 0)
|
|
except Exception:
|
|
raise EvoRuntimeError("MODEL_INPUT_PROJECTION_INVALID") from None
|
|
|
|
|
|
def _callback_messages_payload(messages: Sequence[Sequence[Any]]) -> list[list[Any]]:
|
|
"""Convert LangChain callback messages to canonical-JSON-safe values."""
|
|
|
|
return [
|
|
[
|
|
message_to_dict(message) if isinstance(message, BaseMessage) else message
|
|
for message in batch
|
|
]
|
|
for batch in messages
|
|
]
|
|
|
|
|
|
def _callback_payload_debug_summary(
|
|
payload: Sequence[Sequence[Any]], input_bound_tokens: int
|
|
) -> dict[str, str | int]:
|
|
"""Return an irreversible, content-free summary of a model request."""
|
|
|
|
message_types: dict[str, int] = defaultdict(int)
|
|
content_block_types: dict[str, int] = defaultdict(int)
|
|
tool_calls = 0
|
|
tool_results = 0
|
|
for batch in payload:
|
|
for item in batch:
|
|
if not isinstance(item, Mapping):
|
|
message_types[type(item).__name__] += 1
|
|
continue
|
|
message_type = str(item.get("type") or "unknown")[:64]
|
|
message_types[message_type] += 1
|
|
data = item.get("data")
|
|
data = data if isinstance(data, Mapping) else {}
|
|
if message_type == "tool":
|
|
tool_results += 1
|
|
declared_tool_calls = data.get("tool_calls")
|
|
has_declared_tool_calls = isinstance(declared_tool_calls, list)
|
|
if isinstance(declared_tool_calls, list):
|
|
tool_calls += len(declared_tool_calls)
|
|
content = data.get("content")
|
|
if not isinstance(content, list):
|
|
continue
|
|
for block in content:
|
|
if isinstance(block, Mapping):
|
|
block_type = str(block.get("type") or "unknown")[:64]
|
|
if not has_declared_tool_calls and block_type in {
|
|
"tool_call",
|
|
"tool_use",
|
|
"function_call",
|
|
}:
|
|
tool_calls += 1
|
|
else:
|
|
block_type = type(block).__name__[:64]
|
|
content_block_types[block_type] += 1
|
|
encoded = canonical_json_v1(payload)
|
|
return {
|
|
"input_bound_tokens": input_bound_tokens,
|
|
"input_bytes": len(encoded),
|
|
"input_sha256": hashlib.sha256(encoded).hexdigest()[:24],
|
|
"message_types": ",".join(
|
|
f"{name}:{count}" for name, count in sorted(message_types.items())
|
|
),
|
|
"content_block_types": ",".join(
|
|
f"{name}:{count}" for name, count in sorted(content_block_types.items())
|
|
),
|
|
"tool_calls": tool_calls,
|
|
"tool_results": tool_results,
|
|
}
|
|
|
|
|
|
def _safe_invocation_parameter_value(value: Any, *, nested: bool = False) -> Any:
|
|
"""Project invocation parameters onto a bounded, non-secret debug shape."""
|
|
|
|
if value is None or isinstance(value, (bool, int, float)):
|
|
return value
|
|
if isinstance(value, str):
|
|
return value[:128]
|
|
if isinstance(value, Mapping):
|
|
allowed = (
|
|
_DEBUG_NESTED_PARAMETER_KEYS if nested else _DEBUG_INVOCATION_PARAMETER_KEYS
|
|
)
|
|
return {
|
|
str(key): _safe_invocation_parameter_value(item, nested=True)
|
|
for key, item in sorted(value.items(), key=lambda pair: str(pair[0]))
|
|
if str(key) in allowed
|
|
}
|
|
if isinstance(value, (list, tuple)):
|
|
return [
|
|
_safe_invocation_parameter_value(item, nested=True) for item in value[:16]
|
|
]
|
|
return f"<{type(value).__name__}>"
|
|
|
|
|
|
def _invocation_parameters_debug(plan: InvocationPlan | None) -> str:
|
|
"""Serialize only provider-call parameters that cannot contain credentials."""
|
|
|
|
if plan is None:
|
|
return "{}"
|
|
projected = {
|
|
str(key): _safe_invocation_parameter_value(value, nested=True)
|
|
for key, value in sorted(plan.sdk_params.items(), key=lambda pair: str(pair[0]))
|
|
if str(key) in _DEBUG_INVOCATION_PARAMETER_KEYS
|
|
}
|
|
return json.dumps(
|
|
projected, ensure_ascii=True, sort_keys=True, separators=(",", ":")
|
|
)
|
|
|
|
|
|
def _callback_message_schema_debug(
|
|
payload: Sequence[Sequence[Any]],
|
|
) -> str:
|
|
"""Describe each callback message without logging IDs or message content."""
|
|
|
|
messages: list[dict[str, Any]] = []
|
|
for batch_index, batch in enumerate(payload[:8]):
|
|
for message_index, item in enumerate(batch[:128]):
|
|
if not isinstance(item, Mapping):
|
|
messages.append(
|
|
{
|
|
"batch": batch_index,
|
|
"index": message_index,
|
|
"type": type(item).__name__[:64],
|
|
}
|
|
)
|
|
continue
|
|
data = item.get("data")
|
|
data = data if isinstance(data, Mapping) else {}
|
|
content = data.get("content")
|
|
descriptor: dict[str, Any] = {
|
|
"batch": batch_index,
|
|
"index": message_index,
|
|
"type": str(item.get("type") or "unknown")[:64],
|
|
"content_kind": type(content).__name__,
|
|
"empty": content in (None, "", []),
|
|
}
|
|
if isinstance(content, str):
|
|
descriptor["content_chars"] = len(content)
|
|
elif isinstance(content, list):
|
|
descriptor["content_items"] = len(content)
|
|
descriptor["content_types"] = [
|
|
str(block.get("type") or "dict")[:64]
|
|
if isinstance(block, Mapping)
|
|
else type(block).__name__[:64]
|
|
for block in content[:64]
|
|
]
|
|
descriptor["text_chars"] = sum(
|
|
len(str(block.get("text") or ""))
|
|
for block in content
|
|
if isinstance(block, Mapping)
|
|
)
|
|
additional = data.get("additional_kwargs")
|
|
if isinstance(additional, Mapping) and additional:
|
|
descriptor["additional_keys"] = sorted(
|
|
str(key)[:64] for key in additional
|
|
)
|
|
response_metadata = data.get("response_metadata")
|
|
if isinstance(response_metadata, Mapping) and response_metadata:
|
|
descriptor["response_metadata_keys"] = sorted(
|
|
str(key)[:64] for key in response_metadata
|
|
)
|
|
tool_calls = data.get("tool_calls")
|
|
if isinstance(tool_calls, list) and tool_calls:
|
|
descriptor["tool_calls"] = len(tool_calls)
|
|
if data.get("name") is not None:
|
|
descriptor["has_name"] = True
|
|
if data.get("tool_call_id") is not None:
|
|
descriptor["has_tool_call_id"] = True
|
|
messages.append(descriptor)
|
|
return json.dumps(
|
|
messages, ensure_ascii=True, sort_keys=True, separators=(",", ":")
|
|
)
|
|
|
|
|
|
def _provider_error_message_debug(
|
|
error: BaseException, route: ResolvedRoute | None
|
|
) -> str:
|
|
"""Return a bounded provider explanation with route credentials removed."""
|
|
|
|
body = getattr(error, "body", None)
|
|
candidates: list[Any] = []
|
|
if isinstance(body, Mapping):
|
|
nested = body.get("error")
|
|
if isinstance(nested, Mapping):
|
|
candidates.append(nested.get("message"))
|
|
candidates.append(body.get("message"))
|
|
candidates.append(getattr(error, "message", None))
|
|
candidates.append(str(error))
|
|
message = next(
|
|
(
|
|
str(value)
|
|
for value in candidates
|
|
if isinstance(value, str) and value.strip()
|
|
),
|
|
"",
|
|
)
|
|
if not message:
|
|
return ""
|
|
secrets: list[str] = []
|
|
if route is not None:
|
|
secrets.append(str(route.api_key or ""))
|
|
secrets.extend(str(value) for value in route.default_headers.values())
|
|
for secret in secrets:
|
|
if len(secret) >= 4:
|
|
message = message.replace(secret, "<redacted>")
|
|
message = _AUTH_VALUE_PATTERN.sub(
|
|
lambda match: match.group(1) + "<redacted>", message
|
|
)
|
|
return " ".join(message.split())[:1_024]
|
|
|
|
|
|
def _provider_json_decode_debug(error: BaseException) -> dict[str, str | int]:
|
|
"""Describe a JSON decoder failure without exposing parsed content."""
|
|
|
|
if not isinstance(error, json.JSONDecodeError):
|
|
return {}
|
|
document = error.doc if isinstance(error.doc, str) else ""
|
|
raw_document = document.encode("utf-8", errors="replace")
|
|
position = max(0, min(int(error.pos), len(document)))
|
|
window = document[max(0, position - 128) : min(len(document), position + 129)]
|
|
raw_window = window.encode("utf-8", errors="replace")
|
|
control_counts: dict[int, int] = {}
|
|
for character in window:
|
|
codepoint = ord(character)
|
|
if codepoint < 0x20:
|
|
control_counts[codepoint] = control_counts.get(codepoint, 0) + 1
|
|
control_summary = ",".join(
|
|
f"U+{codepoint:04X}:{count}"
|
|
for codepoint, count in sorted(control_counts.items())[:8]
|
|
)
|
|
event_match = _OPENAI_RESPONSE_EVENT_TYPE_PATTERN.search(document[:4_096])
|
|
frames = traceback.extract_tb(error.__traceback__)[-8:]
|
|
trace = ">".join(
|
|
f"{frame.filename.rsplit('/', 1)[-1]}:{frame.lineno}:{frame.name}"
|
|
for frame in frames
|
|
)[:1_024]
|
|
details: dict[str, str | int] = {
|
|
"provider_json_document_bytes": len(raw_document),
|
|
"provider_json_document_sha256": hashlib.sha256(raw_document).hexdigest()[:24],
|
|
"provider_json_line": int(error.lineno),
|
|
"provider_json_column": int(error.colno),
|
|
"provider_json_position": position,
|
|
"provider_json_position_inside_string": (
|
|
"true" if _json_position_inside_string(document, position) else "false"
|
|
),
|
|
"provider_json_window_bytes": len(raw_window),
|
|
"provider_json_window_sha256": hashlib.sha256(raw_window).hexdigest()[:24],
|
|
"provider_json_trace": trace,
|
|
}
|
|
if position < len(document):
|
|
details["provider_json_invalid_codepoint"] = f"U+{ord(document[position]):04X}"
|
|
if control_summary:
|
|
details["provider_json_control_codepoints"] = control_summary
|
|
if event_match is not None:
|
|
details["provider_json_event_type"] = event_match.group(1)
|
|
return details
|
|
|
|
|
|
def _json_position_inside_string(document: str, position: int) -> bool:
|
|
"""Return lexical string state at a JSON position without parsing its content."""
|
|
|
|
inside_string = False
|
|
escaped = False
|
|
for character in document[: max(0, min(position, len(document)))]:
|
|
if not inside_string:
|
|
if character == '"':
|
|
inside_string = True
|
|
continue
|
|
if escaped:
|
|
escaped = False
|
|
elif character == "\\":
|
|
escaped = True
|
|
elif character == '"':
|
|
inside_string = False
|
|
return inside_string
|
|
|
|
|
|
def _provider_stream_timeout_debug(
|
|
error: BaseException | None,
|
|
) -> dict[str, str | int]:
|
|
"""Extract structured idle-timeout counters without relying on error text."""
|
|
|
|
if error is None:
|
|
return {}
|
|
details: dict[str, str | int] = {}
|
|
chunks_received = getattr(error, "chunks_received", None)
|
|
if (
|
|
isinstance(chunks_received, int)
|
|
and not isinstance(chunks_received, bool)
|
|
and chunks_received >= 0
|
|
):
|
|
details["provider_stream_chunks_received"] = chunks_received
|
|
timeout_seconds = getattr(error, "timeout_s", None)
|
|
if (
|
|
isinstance(timeout_seconds, (int, float))
|
|
and not isinstance(timeout_seconds, bool)
|
|
and math.isfinite(float(timeout_seconds))
|
|
and float(timeout_seconds) >= 0
|
|
):
|
|
details["provider_stream_idle_timeout_ms"] = round(
|
|
float(timeout_seconds) * 1_000
|
|
)
|
|
return details
|
|
|
|
|
|
def _normalize_usage(
|
|
value: Mapping[str, Any] | None,
|
|
) -> Mapping[str, int | str] | None:
|
|
if value is None:
|
|
return None
|
|
try:
|
|
raw_input = value.get("input_tokens", value.get("prompt_tokens"))
|
|
input_details = value.get("input_token_details")
|
|
if not isinstance(input_details, Mapping):
|
|
input_details = value.get("prompt_tokens_details")
|
|
if not isinstance(input_details, Mapping):
|
|
input_details = {}
|
|
raw_cached = value.get(
|
|
"cached_input_tokens",
|
|
value.get(
|
|
"cached_tokens",
|
|
input_details.get(
|
|
"cache_read",
|
|
input_details.get(
|
|
"cached_tokens",
|
|
input_details.get("cache_read_input_tokens"),
|
|
),
|
|
),
|
|
),
|
|
)
|
|
raw_output = value.get("output_tokens", value.get("completion_tokens"))
|
|
input_tokens = int(raw_input) if raw_input is not None else None
|
|
cached_tokens = int(raw_cached) if raw_cached is not None else 0
|
|
output_tokens = int(raw_output) if raw_output is not None else None
|
|
output_details = value.get("output_token_details")
|
|
if not isinstance(output_details, Mapping):
|
|
output_details = value.get("completion_tokens_details")
|
|
if not isinstance(output_details, Mapping):
|
|
output_details = {}
|
|
raw_reasoning = value.get(
|
|
"reasoning_tokens",
|
|
output_details.get("reasoning", output_details.get("reasoning_tokens")),
|
|
)
|
|
reasoning_tokens = int(raw_reasoning) if raw_reasoning is not None else None
|
|
raw_total = value.get("total_tokens")
|
|
total_tokens = int(raw_total) if raw_total is not None else None
|
|
except (TypeError, ValueError):
|
|
return None
|
|
finality = str(value.get("usage_finality") or "confirmed")
|
|
if finality not in {"confirmed", "partial", "unconfirmed"}:
|
|
finality = "unconfirmed"
|
|
request_hash = str(value.get("provider_request_id_hash") or "") or None
|
|
if request_hash is None and value.get("provider_request_id"):
|
|
request_hash = (
|
|
"sha256:"
|
|
+ hashlib.sha256(str(value["provider_request_id"]).encode()).hexdigest()
|
|
)
|
|
return NormalizedUsage(
|
|
input_tokens=input_tokens,
|
|
cached_input_tokens=cached_tokens,
|
|
output_tokens=output_tokens,
|
|
reasoning_tokens=reasoning_tokens,
|
|
total_tokens=total_tokens,
|
|
provider_request_id_hash=request_hash,
|
|
finality=finality, # type: ignore[arg-type]
|
|
).confirmed_projection()
|
|
|
|
|
|
def _actual_cost(usage: Mapping[str, int | str], quote: PricingQuote) -> int:
|
|
input_tokens = int(usage["input_tokens"])
|
|
cached_tokens = int(usage["cached_input_tokens"])
|
|
output_tokens = int(usage["output_tokens"])
|
|
uncached = input_tokens - cached_tokens
|
|
numerator = (
|
|
uncached * quote.input_microunits_per_million
|
|
+ cached_tokens * quote.cached_input_microunits_per_million
|
|
+ output_tokens * quote.output_microunits_per_million
|
|
)
|
|
if numerator > _BIGINT_MAX * quote.unit_scale:
|
|
raise EvoRuntimeError("COST_OVERFLOW")
|
|
return math.ceil(numerator / quote.unit_scale)
|
|
|
|
|
|
def _safe_error_code(
|
|
exc: BaseException, *, fallback: str = "MODEL_PROVIDER_ERROR"
|
|
) -> str:
|
|
if isinstance(exc, EvoRuntimeError):
|
|
return exc.code
|
|
if type(exc).__name__ == "GraphRecursionError":
|
|
return "AGENT_RECURSION_LIMIT_EXCEEDED"
|
|
stable_code = getattr(exc, "code", None)
|
|
if (
|
|
isinstance(stable_code, str)
|
|
and 2 <= len(stable_code) <= 64
|
|
and "A" <= stable_code[0] <= "Z"
|
|
and all(
|
|
character == "_" or character.isdigit() or "A" <= character <= "Z"
|
|
for character in stable_code
|
|
)
|
|
):
|
|
return stable_code
|
|
name = type(exc).__name__.upper()
|
|
if "TIMEOUT" in name:
|
|
return "PROVIDER_TIMEOUT" if fallback == "MODEL_PROVIDER_ERROR" else fallback
|
|
return fallback
|
|
|
|
|
|
def _run_failure_details(
|
|
exc: BaseException, error_code: str
|
|
) -> dict[str, str | int | bool] | None:
|
|
"""Return content-free diagnostics for failures outside the model callback."""
|
|
|
|
if isinstance(exc, EvoRuntimeError) and exc.details:
|
|
merged: dict[str, str | int | bool] = {}
|
|
for item in exc.details:
|
|
for key, value in item.items():
|
|
if isinstance(value, str | int | bool):
|
|
merged[str(key)] = value
|
|
return merged or None
|
|
if error_code != "AGENT_RUNTIME_ERROR":
|
|
return None
|
|
return {
|
|
"failure_stage": "agent_execution",
|
|
"reason": "unclassified_agent_exception",
|
|
"agent_error_type": type(exc).__name__[:128],
|
|
"agent_error_module": type(exc).__module__[:128],
|
|
}
|
|
|
|
|
|
def _provider_failure_details(
|
|
error: BaseException,
|
|
route: ResolvedRoute | None,
|
|
error_code: str,
|
|
) -> dict[str, str | int]:
|
|
"""Project a provider failure into safe, actionable runtime diagnostics."""
|
|
|
|
details: dict[str, str | int] = {
|
|
"failure_stage": "provider_request",
|
|
"reason": _provider_failure_reason(error_code),
|
|
"provider_error_type": type(error).__name__[:128],
|
|
"provider_error_module": type(error).__module__[:128],
|
|
}
|
|
details.update(_provider_json_decode_debug(error))
|
|
details.update(_provider_stream_timeout_debug(error))
|
|
if route is not None:
|
|
plan = route.invocation_plan
|
|
details.update(
|
|
{
|
|
"provider": route.identity.provider_id,
|
|
"model": route.identity.model_id,
|
|
"route_key": route.identity.route_key,
|
|
"endpoint": route.identity.endpoint_name,
|
|
"api_mode": plan.api_mode if plan else route.identity.api_mode,
|
|
"tool_call_transport": (
|
|
plan.tool_call_transport
|
|
if plan
|
|
else route.identity.tool_call_transport
|
|
),
|
|
}
|
|
)
|
|
if plan is not None:
|
|
details["output_token_limit"] = plan.output_token_limit
|
|
details["output_token_parameter"] = plan.output_token_parameter
|
|
details["invocation_plan_hash"] = plan.plan_hash
|
|
elif route.max_output_tokens > 0:
|
|
details["output_token_limit"] = route.max_output_tokens
|
|
if status_code := _provider_http_status(error):
|
|
details["http_status"] = status_code
|
|
if provider_code := _provider_error_code(error):
|
|
details["provider_error_code"] = provider_code
|
|
if provider_parameter := _provider_error_parameter(error):
|
|
details["provider_error_parameter"] = provider_parameter
|
|
return details
|
|
|
|
|
|
def _provider_http_status(error: BaseException) -> int | None:
|
|
status = getattr(error, "status_code", None)
|
|
if status is None:
|
|
status = getattr(getattr(error, "response", None), "status_code", None)
|
|
return status if isinstance(status, int) and 100 <= status <= 599 else None
|
|
|
|
|
|
def _provider_error_code(error: BaseException) -> str | None:
|
|
"""Extract only a token-like provider code; never expose provider messages."""
|
|
|
|
values = [getattr(error, "code", None)]
|
|
body = getattr(error, "body", None)
|
|
if isinstance(body, Mapping):
|
|
values.extend((body.get("code"), body.get("type")))
|
|
nested = body.get("error")
|
|
if isinstance(nested, Mapping):
|
|
values.extend((nested.get("code"), nested.get("type")))
|
|
for value in values:
|
|
if not isinstance(value, str):
|
|
continue
|
|
normalized = value.strip()
|
|
if 1 <= len(normalized) <= 128 and all(
|
|
character.isascii() and (character.isalnum() or character in "_.:-")
|
|
for character in normalized
|
|
):
|
|
return normalized
|
|
return None
|
|
|
|
|
|
def _provider_error_parameter(error: BaseException) -> str | None:
|
|
"""Return only a known rejected request field, never provider prose."""
|
|
|
|
body = getattr(error, "body", None)
|
|
mappings = [body] if isinstance(body, Mapping) else []
|
|
if isinstance(body, Mapping) and isinstance(body.get("error"), Mapping):
|
|
mappings.append(body["error"])
|
|
for mapping in mappings:
|
|
for key in ("param", "parameter", "field"):
|
|
value = mapping.get(key)
|
|
if isinstance(value, str):
|
|
normalized = value.strip()
|
|
if normalized in _KNOWN_PROVIDER_REQUEST_FIELDS:
|
|
return normalized
|
|
message = mapping.get("message")
|
|
if isinstance(message, str):
|
|
lowered = message.lower()
|
|
for parameter in _KNOWN_PROVIDER_REQUEST_FIELDS:
|
|
if parameter in lowered:
|
|
return parameter
|
|
return None
|
|
|
|
|
|
_KNOWN_PROVIDER_REQUEST_FIELDS = (
|
|
"max_output_tokens",
|
|
"max_completion_tokens",
|
|
"max_tokens",
|
|
"reasoning_effort",
|
|
"reasoning",
|
|
"temperature",
|
|
"top_p",
|
|
"tool_choice",
|
|
"tools",
|
|
"response_format",
|
|
"stream",
|
|
"messages",
|
|
)
|
|
|
|
|
|
def _provider_failure_reason(error_code: str) -> str:
|
|
return {
|
|
"MODEL_PROVIDER_REQUEST_REJECTED": "provider_rejected_request",
|
|
"MODEL_AUTHENTICATION_FAILED": "provider_authentication_failed",
|
|
"MODEL_NOT_FOUND": "provider_model_or_endpoint_not_found",
|
|
"MODEL_RATE_LIMITED": "provider_rate_limited",
|
|
"MODEL_TIMEOUT": "provider_timeout",
|
|
}.get(error_code, "provider_request_failed")
|
|
|
|
|
|
def _title_prompt(source_text: str) -> str:
|
|
normalized = " ".join(source_text.split())[:500]
|
|
return f"Generate a concise title of no more than 30 characters. Return only the title.\n\n{normalized}"
|
|
|
|
|
|
def _response_text(response: Any) -> str:
|
|
content = getattr(response, "content", response)
|
|
if isinstance(content, str):
|
|
return content
|
|
if isinstance(content, list):
|
|
return "".join(
|
|
str(item.get("text") or "") if isinstance(item, Mapping) else str(item)
|
|
for item in content
|
|
)
|
|
return str(content)
|
|
|
|
|
|
def _usage_from_response(response: Any) -> Mapping[str, Any] | None:
|
|
usage = getattr(response, "usage_metadata", None)
|
|
if isinstance(usage, Mapping):
|
|
return usage
|
|
llm_output = getattr(response, "llm_output", None)
|
|
if isinstance(llm_output, Mapping):
|
|
candidate = llm_output.get("usage") or llm_output.get("token_usage")
|
|
if isinstance(candidate, Mapping):
|
|
return candidate
|
|
generations = getattr(response, "generations", None)
|
|
if isinstance(generations, Sequence) and not isinstance(generations, str | bytes):
|
|
for generation_group in generations:
|
|
items = (
|
|
generation_group
|
|
if isinstance(generation_group, Sequence)
|
|
and not isinstance(generation_group, str | bytes)
|
|
else (generation_group,)
|
|
)
|
|
for generation in items:
|
|
message = getattr(generation, "message", None)
|
|
for value in (generation, message):
|
|
nested_usage = getattr(value, "usage_metadata", None)
|
|
if isinstance(nested_usage, Mapping):
|
|
return nested_usage
|
|
response_metadata = getattr(value, "response_metadata", None)
|
|
if isinstance(response_metadata, Mapping):
|
|
candidate = response_metadata.get(
|
|
"usage"
|
|
) or response_metadata.get("token_usage")
|
|
if isinstance(candidate, Mapping):
|
|
return candidate
|
|
generation_info = getattr(generation, "generation_info", None)
|
|
if isinstance(generation_info, Mapping):
|
|
candidate = generation_info.get("usage") or generation_info.get(
|
|
"token_usage"
|
|
)
|
|
if isinstance(candidate, Mapping):
|
|
return candidate
|
|
return None
|