Files

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