refactor(agent/transports): shared codex turn driver, build_kwargs phase helpers, dispatch tables in projector/accumulator, base finish-reason mapping

This commit is contained in:
Teknium
2026-09-02 13:18:33 -07:00
parent 00cd12f7e4
commit d2b3e545cb
11 changed files with 1397 additions and 3410 deletions

View File

@@ -1,21 +1,20 @@
"""Transport layer types and registry for provider response normalization.
"""Transport registry for provider response normalization.
Usage:
from agent.transports import get_transport
transport = get_transport("anthropic_messages")
result = transport.normalize_response(raw_response)
"""
from agent.transports.types import (
from agent.transports.types import ( # noqa: F401
NormalizedResponse,
ToolCall,
Usage,
build_tool_call,
map_finish_reason,
) # noqa: F401
)
_REGISTRY: dict = {}
_discovered: bool = False
_TRANSPORT_MODULES = ("anthropic", "codex", "chat_completions", "bedrock")
def register_transport(api_mode: str, transport_cls: type) -> None:
@@ -24,45 +23,26 @@ def register_transport(api_mode: str, transport_cls: type) -> None:
def get_transport(api_mode: str):
"""Get a transport instance for the given api_mode.
Returns None if no transport is registered for this api_mode.
This allows gradual migration — call sites can check for None
and fall back to the legacy code path.
"""
global _discovered
"""Return a transport instance for ``api_mode``, or None so callers can fall back to the legacy path."""
if not _discovered:
_discover_transports()
cls = _REGISTRY.get(api_mode)
if cls is None:
# The registry can be partially populated when a specific transport
# module was imported directly (for example chat_completions before
# codex). Discover on misses, not only when the registry is empty, so
# test/order-dependent imports do not make valid api_modes unavailable.
# A directly-imported transport module leaves the registry partially
# populated; discover on misses so import order can't hide a valid api_mode.
_discover_transports()
cls = _REGISTRY.get(api_mode)
if cls is None:
return None
return cls()
return None if cls is None else cls()
def _discover_transports() -> None:
"""Import all transport modules to trigger auto-registration."""
global _discovered
_discovered = True
try:
import agent.transports.anthropic # noqa: F401
except ImportError:
pass
try:
import agent.transports.codex # noqa: F401
except ImportError:
pass
try:
import agent.transports.chat_completions # noqa: F401
except ImportError:
pass
try:
import agent.transports.bedrock # noqa: F401
except ImportError:
pass
import importlib
for name in _TRANSPORT_MODULES:
try:
importlib.import_module(f"agent.transports.{name}")
except ImportError:
pass

View File

@@ -1,36 +1,57 @@
"""Anthropic Messages API transport.
Delegates to the existing adapter functions in agent/anthropic_adapter.py.
This transport owns format conversion and normalization — NOT client lifecycle.
Delegates format conversion to agent/anthropic_adapter.py; owns normalization,
not client lifecycle.
"""
from typing import Any, Dict, List, Optional
from agent.transports.base import ProviderTransport
from agent.transports.types import NormalizedResponse
from agent.transports.types import NormalizedResponse, ToolCall
_MCP_PREFIX = "mcp__"
def _unprefix_oauth_tool_name(name: str) -> str:
"""Reverse the OAuth-wire ``mcp__`` prefix back to the registered tool name.
Two originals map onto one wire name (``mcp__read_file`` <- ``read_file``;
``mcp__linear_get_issue`` <- ``mcp_linear_get_issue``), so resolve by registry
lookup, never rewriting a name that already resolves natively (GH-25255).
OAuth wire aliases (e.g. chat_history_lookup -> session_search) are checked
LAST so a real tool registered under the wire name still wins.
"""
from agent.anthropic_adapter import _OAUTH_TOOL_NAME_REVERSE_ALIASES
from tools.registry import registry as _tool_registry
bare = name[len(_MCP_PREFIX):]
for candidate in (name, "mcp_" + bare, bare):
if _tool_registry.get_entry(candidate):
return candidate
return _OAUTH_TOOL_NAME_REVERSE_ALIASES.get(bare, name)
class AnthropicTransport(ProviderTransport):
"""Transport for api_mode='anthropic_messages'.
"""Transport for api_mode='anthropic_messages'."""
Wraps the existing functions in anthropic_adapter.py behind the
ProviderTransport ABC. Each method delegates — no logic is duplicated.
"""
_STOP_REASON_MAP = {
"end_turn": "stop",
"tool_use": "tool_calls",
"max_tokens": "length",
"stop_sequence": "stop",
"refusal": "content_filter",
"model_context_window_exceeded": "length",
}
@property
def api_mode(self) -> str:
return "anthropic_messages"
def convert_messages(self, messages: List[Dict[str, Any]], **kwargs) -> Any:
"""Convert OpenAI messages to Anthropic (system, messages) tuple.
kwargs:
base_url: Optional[str] — affects thinking signature handling.
"""
"""Convert OpenAI messages to an Anthropic (system, messages) tuple; ``base_url`` affects thinking-signature handling."""
from agent.anthropic_adapter import convert_messages_to_anthropic
base_url = kwargs.get("base_url")
return convert_messages_to_anthropic(messages, base_url=base_url)
return convert_messages_to_anthropic(messages, base_url=kwargs.get("base_url"))
def convert_tools(self, tools: List[Dict[str, Any]]) -> Any:
"""Convert OpenAI tool schemas to Anthropic input_schema format."""
@@ -45,21 +66,7 @@ class AnthropicTransport(ProviderTransport):
tools: Optional[List[Dict[str, Any]]] = None,
**params,
) -> Dict[str, Any]:
"""Build Anthropic messages.create() kwargs.
Calls convert_messages and convert_tools internally.
params (all optional):
max_tokens: int
reasoning_config: dict | None
tool_choice: str | None
is_oauth: bool
preserve_dots: bool
context_length: int | None
base_url: str | None
fast_mode: bool
drop_context_1m_beta: bool
"""
"""Build Anthropic messages.create() kwargs (converts messages and tools internally)."""
from agent.anthropic_adapter import build_anthropic_kwargs
return build_anthropic_kwargs(
@@ -78,46 +85,24 @@ class AnthropicTransport(ProviderTransport):
)
def normalize_response(self, response: Any, **kwargs) -> NormalizedResponse:
"""Normalize Anthropic response to NormalizedResponse.
Parses content blocks (text, thinking, tool_use), maps stop_reason
to OpenAI finish_reason, and collects reasoning_details in provider_data.
"""
"""Parse content blocks (text/thinking/tool_use), map stop_reason, collect reasoning_details."""
import json
from agent.anthropic_adapter import (
_OAUTH_TOOL_NAME_REVERSE_ALIASES,
_sanitize_replay_block,
_to_plain_data,
)
from agent.transports.types import ToolCall
from agent.anthropic_adapter import _sanitize_replay_block, _to_plain_data
strip_tool_prefix = kwargs.get("strip_tool_prefix", False)
_MCP_PREFIX = "mcp__"
text_parts = []
reasoning_parts = []
reasoning_details = []
tool_calls = []
# Verbatim, order-preserving copy of every content block in the turn.
# Anthropic signs each thinking block against the turn content that
# PRECEDES it at its position; when a turn interleaves thinking and
# tool_use (adaptive/interleaved thinking, Claude 4.6+), the parallel
# reasoning_details + tool_calls lists below lose that cross-type
# ordering. Replaying the latest assistant message in the wrong order
# invalidates the signatures -> HTTP 400 "thinking ... blocks in the
# latest assistant message cannot be modified". Preserve the exact
# block sequence here so the adapter can replay it unchanged. See
# tests/agent/test_anthropic_thinking_block_order.py.
text_parts, reasoning_parts, reasoning_details, tool_calls = [], [], [], []
# Anthropic signs each thinking block against the blocks that PRECEDE it.
# When thinking interleaves with tool_use, the parallel reasoning_details +
# tool_calls lists lose that ordering and replay -> HTTP 400 "thinking ...
# blocks cannot be modified". Keep the exact sequence for the adapter.
ordered_blocks = []
for block in response.content:
block_dict = _to_plain_data(block)
clean_block = None
if isinstance(block_dict, dict):
# Sanitize at capture so output-only SDK fields (parsed_output,
# caller, citations=None, …) never persist to state.db and leak
# back as request input on replay → HTTP 400 "Extra inputs are
# not permitted". Defence-in-depth with the replay-side sanitize.
# Sanitize at capture so output-only SDK fields never persist to
# state.db and leak back as request input on replay (HTTP 400).
clean_block = _sanitize_replay_block(block_dict)
if clean_block is not None:
ordered_blocks.append(clean_block)
@@ -126,9 +111,7 @@ class AnthropicTransport(ProviderTransport):
elif block.type in ("thinking", "redacted_thinking"):
if block.type == "thinking":
reasoning_parts.append(block.thinking)
# Use the sanitized block (clean_block) for reasoning_details too,
# since _extract_preserved_thinking_blocks replays these on the
# non-ordered path. Falls back to raw only if sanitize dropped it.
# Prefer the sanitized block (replayed on the non-ordered path); raw only if sanitize dropped it.
if isinstance(clean_block, dict):
reasoning_details.append(clean_block)
elif isinstance(block_dict, dict):
@@ -136,126 +119,49 @@ class AnthropicTransport(ProviderTransport):
elif block.type == "tool_use":
name = block.name
if strip_tool_prefix and name.startswith(_MCP_PREFIX):
# On the OAuth wire every tool carries a double-underscore
# ``mcp__`` prefix (added in build_anthropic_kwargs to avoid
# Anthropic's single-underscore third-party classifier).
# Reverse it back to the name the registry/dispatcher knows.
# Two original forms map onto the same ``mcp__`` wire name:
# ``mcp__read_file`` <- bare native tool ``read_file``
# ``mcp__linear_get_issue`` <- MCP server tool
# ``mcp_linear_get_issue``
# Resolve by registry lookup, preferring whichever original
# is actually registered; never rewrite a name the LLM used
# that already resolves natively. GH-25255.
from tools.registry import registry as _tool_registry
if not _tool_registry.get_entry(name):
bare = name[len(_MCP_PREFIX):] # read_file
single = "mcp_" + bare # mcp_read_file / mcp_linear_get_issue
if _tool_registry.get_entry(single):
name = single
elif _tool_registry.get_entry(bare):
name = bare
elif bare in _OAUTH_TOOL_NAME_REVERSE_ALIASES:
# OAuth wire alias (e.g. chat_history_lookup ->
# session_search, #65365). Checked LAST so a real
# tool actually registered under the wire name
# still wins — same GH-25255 precedence.
name = _OAUTH_TOOL_NAME_REVERSE_ALIASES[bare]
tool_calls.append(
ToolCall(
id=block.id,
name=name,
arguments=json.dumps(block.input),
)
)
finish_reason = self._STOP_REASON_MAP.get(response.stop_reason, "stop")
name = _unprefix_oauth_tool_name(name)
tool_calls.append(ToolCall(id=block.id, name=name, arguments=json.dumps(block.input)))
provider_data = {}
if reasoning_details:
provider_data["reasoning_details"] = reasoning_details
# Only worth carrying the ordered-blocks channel when the turn
# actually interleaves signed thinking with tool_use — that's the
# only shape the parallel lists reconstruct incorrectly. A turn that
# is purely text, or thinking-then-tools with a single leading
# thinking block, replays correctly without it.
# Carry the ordered channel only for the one shape the parallel lists
# reconstruct wrongly: signed thinking interleaved with tool_use.
_has_signed_thinking = any(
isinstance(b, dict)
and b.get("type") in ("thinking", "redacted_thinking")
and (b.get("signature") or b.get("data"))
isinstance(b, dict) and b.get("type") in ("thinking", "redacted_thinking") and (b.get("signature") or b.get("data"))
for b in ordered_blocks
)
_has_tool_use = any(
isinstance(b, dict) and b.get("type") == "tool_use"
for b in ordered_blocks
)
if _has_signed_thinking and _has_tool_use:
if _has_signed_thinking and any(isinstance(b, dict) and b.get("type") == "tool_use" for b in ordered_blocks):
provider_data["anthropic_content_blocks"] = ordered_blocks
return NormalizedResponse(
content="\n".join(text_parts) if text_parts else None,
tool_calls=tool_calls or None,
finish_reason=finish_reason,
finish_reason=self.map_finish_reason(response.stop_reason),
reasoning="\n\n".join(reasoning_parts) if reasoning_parts else None,
usage=None,
provider_data=provider_data or None,
)
def validate_response(self, response: Any) -> bool:
"""Check Anthropic response structure is valid.
An empty content list is legitimate for terminal stop reasons that
carry no text payload:
- ``end_turn`` — the model's canonical "nothing more to add" after a
tool turn that already delivered the user-facing text.
- ``refusal`` — the model declined to respond (Claude 4.5+). The
Messages API returns an empty ``content`` list with this stop
reason. Treating it as invalid sends a deterministic refusal into
the invalid-response retry loop, which reproduces the refusal on
every attempt and surfaces a misleading "rate limited / invalid
response" error instead of the refusal. ``normalize_response`` maps
``refusal`` → ``content_filter`` so the agent loop's refusal handler
can surface it.
Treating either as invalid falsely retries a completed response.
"""
if response is None:
return False
content_blocks = getattr(response, "content", None)
"""Structural check. An empty content list is legitimate for ``end_turn`` (nothing to add
after a tool turn) and ``refusal`` (Claude 4.5+ declines with empty content); treating
either as invalid would retry a completed/deterministic response forever."""
content_blocks = getattr(response, "content", None) if response is not None else None
if not isinstance(content_blocks, list):
return False
if not content_blocks:
return getattr(response, "stop_reason", None) in {"end_turn", "refusal"}
return True
return bool(content_blocks) or getattr(response, "stop_reason", None) in {"end_turn", "refusal"}
def extract_cache_stats(self, response: Any) -> Optional[Dict[str, int]]:
"""Extract Anthropic cache_read and cache_creation token counts."""
"""Anthropic cache_read / cache_creation token counts."""
usage = getattr(response, "usage", None)
if usage is None:
return None
cached = getattr(usage, "cache_read_input_tokens", 0) or 0
written = getattr(usage, "cache_creation_input_tokens", 0) or 0
if cached or written:
return {"cached_tokens": cached, "creation_tokens": written}
return None
# Promote the adapter's canonical mapping to module level so it's shared
_STOP_REASON_MAP = {
"end_turn": "stop",
"tool_use": "tool_calls",
"max_tokens": "length",
"stop_sequence": "stop",
"refusal": "content_filter",
"model_context_window_exceeded": "length",
}
def map_finish_reason(self, raw_reason: str) -> str:
"""Map Anthropic stop_reason to OpenAI finish_reason."""
return self._STOP_REASON_MAP.get(raw_reason, "stop")
return {"cached_tokens": cached, "creation_tokens": written} if cached or written else None
# Auto-register on import
from agent.transports import register_transport # noqa: E402
register_transport("anthropic_messages", AnthropicTransport)

View File

@@ -1,10 +1,9 @@
"""Abstract base for provider transports.
A transport owns the data path for one api_mode:
convert_messages → convert_tools → build_kwargs → normalize_response
It does NOT own: client construction, streaming, credential refresh,
prompt caching, interrupt handling, or retry logic. Those stay on AIAgent.
convert_messages -> convert_tools -> build_kwargs -> normalize_response
It does NOT own client construction, streaming, credential refresh, prompt
caching, interrupt handling, or retry logic — those stay on AIAgent.
"""
from abc import ABC, abstractmethod
@@ -16,28 +15,22 @@ from agent.transports.types import NormalizedResponse
class ProviderTransport(ABC):
"""Base class for provider-specific format conversion and normalization."""
# Provider stop_reason -> OpenAI finish_reason. ``None`` means the provider
# already speaks OpenAI vocabulary and map_finish_reason passes through.
_STOP_REASON_MAP: Optional[Dict[str, str]] = None
@property
@abstractmethod
def api_mode(self) -> str:
"""The api_mode string this transport handles (e.g. 'anthropic_messages')."""
...
@abstractmethod
def convert_messages(self, messages: List[Dict[str, Any]], **kwargs) -> Any:
"""Convert OpenAI-format messages to provider-native format.
Returns provider-specific structure (e.g. (system, messages) for Anthropic,
or the messages list unchanged for chat_completions).
"""
...
"""Convert OpenAI-format messages to the provider-native structure (e.g. (system, messages) for Anthropic)."""
@abstractmethod
def convert_tools(self, tools: List[Dict[str, Any]]) -> Any:
"""Convert OpenAI-format tool definitions to provider-native format.
Returns provider-specific tool list (e.g. Anthropic input_schema format).
"""
...
"""Convert OpenAI-format tool definitions to provider-native format."""
@abstractmethod
def build_kwargs(
@@ -47,43 +40,20 @@ class ProviderTransport(ABC):
tools: Optional[List[Dict[str, Any]]] = None,
**params,
) -> Dict[str, Any]:
"""Build the complete API call kwargs dict.
This is the primary entry point — it typically calls convert_messages()
and convert_tools() internally, then adds model-specific config.
Returns a dict ready to be passed to the provider's SDK client.
"""
...
"""Primary entry point: convert messages/tools and return kwargs ready for the provider SDK."""
@abstractmethod
def normalize_response(self, response: Any, **kwargs) -> NormalizedResponse:
"""Normalize a raw provider response to the shared NormalizedResponse type.
This is the only method that returns a transport-layer type.
"""
...
"""Normalize a raw provider response to NormalizedResponse (the only transport-layer return type)."""
def validate_response(self, response: Any) -> bool:
"""Optional: check if the raw response is structurally valid.
Returns True if valid, False if the response should be treated as invalid.
Default implementation always returns True.
"""
"""Optional structural validity check; default accepts everything."""
return True
def extract_cache_stats(self, response: Any) -> Optional[Dict[str, int]]:
"""Optional: extract provider-specific cache hit/creation stats.
Returns dict with 'cached_tokens' and 'creation_tokens', or None.
Default returns None.
"""
"""Optional: ``{'cached_tokens', 'creation_tokens'}`` or None (default)."""
return None
def map_finish_reason(self, raw_reason: str) -> str:
"""Optional: map provider-specific stop reason to OpenAI equivalent.
Default returns the raw reason unchanged. Override for providers
with different stop reason vocabularies.
"""
return raw_reason
"""Map a provider stop reason via ``_STOP_REASON_MAP`` (unknown -> 'stop'); passthrough when no map."""
return raw_reason if self._STOP_REASON_MAP is None else self._STOP_REASON_MAP.get(raw_reason, "stop")

View File

@@ -1,9 +1,7 @@
"""AWS Bedrock Converse API transport.
Delegates to the existing adapter functions in agent/bedrock_adapter.py.
Bedrock uses its own boto3 client (not the OpenAI SDK), so the transport
owns format conversion and normalization, while client construction and
boto3 calls stay on AIAgent.
Delegates format conversion to agent/bedrock_adapter.py. Bedrock uses its own
boto3 client, so client construction and calls stay on AIAgent.
"""
from typing import Any, Dict, List, Optional
@@ -15,6 +13,16 @@ from agent.transports.types import NormalizedResponse, ToolCall, Usage
class BedrockTransport(ProviderTransport):
"""Transport for api_mode='bedrock_converse'."""
# The adapter already maps inside normalize_converse_response; this serves raw-response access.
_STOP_REASON_MAP = {
"end_turn": "stop",
"tool_use": "tool_calls",
"max_tokens": "length",
"stop_sequence": "stop",
"guardrail_intervened": "content_filter",
"content_filtered": "content_filter",
}
@property
def api_mode(self) -> str:
return "bedrock_converse"
@@ -36,76 +44,36 @@ class BedrockTransport(ProviderTransport):
tools: Optional[List[Dict[str, Any]]] = None,
**params,
) -> Dict[str, Any]:
"""Build Bedrock converse() kwargs.
Calls convert_messages and convert_tools internally.
params:
max_tokens: int — output token limit (default 4096)
temperature: float | None
guardrail_config: dict | None — Bedrock guardrails
region: str — AWS region (default 'us-east-1')
"""
"""Build converse() kwargs; params: max_tokens (4096), temperature, guardrail_config, region ('us-east-1')."""
from agent.bedrock_adapter import build_converse_kwargs
region = params.get("region", "us-east-1")
guardrail = params.get("guardrail_config")
kwargs = build_converse_kwargs(
model=model,
messages=messages,
tools=tools,
max_tokens=params.get("max_tokens", 4096),
temperature=params.get("temperature"),
guardrail_config=guardrail,
guardrail_config=params.get("guardrail_config"),
)
# Sentinel keys for dispatch — agent pops these before the boto3 call
kwargs["__bedrock_converse__"] = True
kwargs["__bedrock_region__"] = region
kwargs["__bedrock_region__"] = params.get("region", "us-east-1")
return kwargs
def normalize_response(self, response: Any, **kwargs) -> NormalizedResponse:
"""Normalize Bedrock response to NormalizedResponse.
Handles two shapes:
1. Raw boto3 dict (from direct converse() calls)
2. Already-normalized SimpleNamespace with .choices (from dispatch site)
"""
"""Normalize either a raw boto3 dict or an already-normalized SimpleNamespace with .choices."""
from agent.bedrock_adapter import normalize_converse_response
# Normalize to OpenAI-compatible SimpleNamespace
if hasattr(response, "choices") and response.choices:
# Already normalized at dispatch site
ns = response
else:
# Raw boto3 dict
ns = normalize_converse_response(response)
ns = response if hasattr(response, "choices") and response.choices else normalize_converse_response(response)
choice = ns.choices[0]
msg = choice.message
finish_reason = choice.finish_reason or "stop"
tool_calls = None
if msg.tool_calls:
tool_calls = [
ToolCall(
id=tc.id,
name=tc.function.name,
arguments=tc.function.arguments,
)
for tc in msg.tool_calls
]
tool_calls = [ToolCall(id=tc.id, name=tc.function.name, arguments=tc.function.arguments) for tc in msg.tool_calls]
usage = None
if hasattr(ns, "usage") and ns.usage:
u = ns.usage
usage = Usage(
prompt_tokens=getattr(u, "prompt_tokens", 0) or 0,
completion_tokens=getattr(u, "completion_tokens", 0) or 0,
total_tokens=getattr(u, "total_tokens", 0) or 0,
)
reasoning = getattr(msg, "reasoning", None) or getattr(msg, "reasoning_content", None)
usage = Usage.from_openai(ns.usage)
provider_data = {}
if getattr(msg, "reasoning_details", None):
@@ -116,46 +84,19 @@ class BedrockTransport(ProviderTransport):
return NormalizedResponse(
content=msg.content,
tool_calls=tool_calls,
finish_reason=finish_reason,
reasoning=reasoning,
finish_reason=choice.finish_reason or "stop",
reasoning=getattr(msg, "reasoning", None) or getattr(msg, "reasoning_content", None),
usage=usage,
provider_data=provider_data or None,
)
def validate_response(self, response: Any) -> bool:
"""Check Bedrock response structure.
After normalize_converse_response, the response has OpenAI-compatible
.choices — same check as chat_completions.
"""
if response is None:
return False
# Raw Bedrock dict response — check for 'output' key
"""Raw Bedrock dict needs an 'output' key; a normalized namespace needs non-empty .choices."""
if isinstance(response, dict):
return "output" in response
# Already-normalized SimpleNamespace
if hasattr(response, "choices"):
return bool(response.choices)
return False
def map_finish_reason(self, raw_reason: str) -> str:
"""Map Bedrock stop reason to OpenAI finish_reason.
The adapter already does this mapping inside normalize_converse_response,
so this is only used for direct access to raw responses.
"""
_MAP = {
"end_turn": "stop",
"tool_use": "tool_calls",
"max_tokens": "length",
"stop_sequence": "stop",
"guardrail_intervened": "content_filter",
"content_filtered": "content_filter",
}
return _MAP.get(raw_reason, "stop")
return bool(getattr(response, "choices", None)) if response is not None else False
# Auto-register on import
from agent.transports import register_transport # noqa: E402
register_transport("bedrock_converse", BedrockTransport)

File diff suppressed because it is too large Load Diff

File diff suppressed because it is too large Load Diff

View File

@@ -1,17 +1,10 @@
"""Codex app-server JSON-RPC client.
Speaks the protocol documented in codex-rs/app-server/README.md (codex 0.125+).
Transport is newline-delimited JSON-RPC 2.0 over stdio: spawn `codex app-server`,
do an `initialize` handshake, then drive `thread/start` + `turn/start` and
consume streaming `item/*` notifications until `turn/completed`.
This module is the wire-level speaker only. Higher-level concerns (event
projection into Hermes' display, approval bridging, transcript projection into
AIAgent.messages, plugin migration) live in sibling modules.
Status: optional opt-in runtime gated behind `model.openai_runtime ==
"codex_app_server"`. Hermes' default tool dispatch is unchanged when this
runtime is not selected.
Newline-delimited JSON-RPC 2.0 over stdio to ``codex app-server`` (codex 0.125+):
``initialize`` handshake, then ``thread/start`` + ``turn/start`` with streaming
``item/*`` notifications until ``turn/completed``. Wire-level speaker only —
projection, approvals and transcript handling live in sibling modules.
Opt-in runtime gated behind ``model.openai_runtime == "codex_app_server"``.
"""
from __future__ import annotations
@@ -19,6 +12,7 @@ from __future__ import annotations
import json
import os
import queue
import re
import subprocess
import threading
import time
@@ -27,8 +21,6 @@ from typing import Any, Optional
from tools.environments.local import hermes_subprocess_env
# Default minimum codex version we test against. The PR sets this from the
# `codex --version` parsed at install time; bumping is a one-line change here.
MIN_CODEX_VERSION = (0, 125, 0)
@@ -52,20 +44,12 @@ class _Pending:
class CodexAppServerClient:
"""Minimal JSON-RPC 2.0 client for `codex app-server` over stdio.
"""Minimal synchronous JSON-RPC 2.0 client for ``codex app-server`` over stdio.
Threading model:
- Spawning thread (caller) drives request/response pairs synchronously.
- One reader thread parses stdout, dispatches replies to the right
pending future, and routes notifications + server-initiated requests
to bounded queues that the caller drains on their own cadence.
- One reader thread captures stderr for diagnostics; codex emits
tracing logs there at RUST_LOG-controlled levels.
Intentionally NOT async. AIAgent.run_conversation() is synchronous and
runs on the main thread; layering asyncio just to drive a stdio child
creates surprising interrupt semantics. We use blocking queues with
timeouts and rely on `turn/interrupt` for cancellation.
The caller drives request/response pairs; one reader thread routes replies
to pending queues and notifications / server requests to bounded queues;
another captures stderr. Intentionally NOT async: AIAgent.run_conversation()
is synchronous and cancellation goes through ``turn/interrupt``.
"""
def __init__(
@@ -76,17 +60,8 @@ class CodexAppServerClient:
env: Optional[dict[str, str]] = None,
) -> None:
self._codex_bin = codex_bin
# codex app-server is a model-driving CLI executor: it runs a
# model-chosen agentic loop that executes shell commands, so it
# legitimately needs LLM provider credentials (inherit_credentials=True)
# to authenticate against the model endpoint. But the previous
# `os.environ.copy()` also handed it every Tier-1 Hermes secret — gateway
# bot tokens, GitHub auth, Modal/Daytona infra tokens, the dashboard
# session token, AUXILIARY_* side-LLM keys, GATEWAY_RELAY_* auth — none
# of which a coding subprocess has any use for. Route through the
# centralized helper so Tier-1 + dynamic-internal secrets are always
# stripped while provider creds still flow, matching copilot_acp_client
# (#29157 sibling spawn-site gap).
# codex needs LLM provider creds (inherit_credentials=True) but must not
# receive Tier-1 Hermes secrets (gateway/GitHub/infra tokens) — #29157.
spawn_env = hermes_subprocess_env(inherit_credentials=True)
if env:
spawn_env.update(env)
@@ -94,51 +69,28 @@ class CodexAppServerClient:
spawn_env["CODEX_HOME"] = codex_home
app_server_args = list(extra_args or [])
# Kanban workers must be able to write their handoff/status back to
# the board DB, which lives outside the per-task workspace. Keep the
# Codex sandbox on, but add the Kanban root as the only extra writable
# root. Without this, codex-runtime workers finish their actual work
# but crash/block when kanban_complete/kanban_block writes SQLite.
# Kanban workers must write handoff/status to the board DB outside the
# workspace: keep the sandbox on, add the Kanban root as writable.
if spawn_env.get("HERMES_KANBAN_TASK"):
kanban_db = spawn_env.get("HERMES_KANBAN_DB")
kanban_root = (
os.path.dirname(kanban_db)
if kanban_db
else spawn_env.get(
"HERMES_KANBAN_ROOT",
os.path.join(
spawn_env.get("HERMES_HOME", os.path.expanduser("~/.hermes")),
"kanban",
),
)
)
app_server_args.extend(
[
"-c",
'sandbox_mode="workspace-write"',
"-c",
f'sandbox_workspace_write.writable_roots=["{kanban_root}"]',
"-c",
"sandbox_workspace_write.network_access=false",
]
)
default_root = os.path.join(spawn_env.get("HERMES_HOME", os.path.expanduser("~/.hermes")), "kanban")
kanban_root = os.path.dirname(kanban_db) if kanban_db else spawn_env.get("HERMES_KANBAN_ROOT", default_root)
app_server_args.extend([
"-c", 'sandbox_mode="workspace-write"',
"-c", f'sandbox_workspace_write.writable_roots=["{kanban_root}"]',
"-c", "sandbox_workspace_write.network_access=false",
])
cmd = [codex_bin, "app-server"] + app_server_args
# Codex emits tracing to stderr; default WARN keeps it quiet for users.
spawn_env.setdefault("RUST_LOG", "warn")
# Hide the console the codex child would otherwise flash on Windows
# (#56747). Hide-only — stdio pipes stay intact for the app-server wire.
# Hide the console flash on Windows (#56747); stdio pipes stay intact.
from hermes_cli._subprocess_compat import windows_hide_flags
self._proc = subprocess.Popen(
cmd,
stdin=subprocess.PIPE,
stdout=subprocess.PIPE,
stderr=subprocess.PIPE,
bufsize=0,
env=spawn_env,
creationflags=windows_hide_flags(),
cmd, stdin=subprocess.PIPE, stdout=subprocess.PIPE, stderr=subprocess.PIPE,
bufsize=0, env=spawn_env, creationflags=windows_hide_flags(),
)
self._next_id = 1
self._pending: dict[int, _Pending] = {}
@@ -151,12 +103,10 @@ class CodexAppServerClient:
self._initialized = False
self._reader = threading.Thread(target=self._read_stdout, daemon=True)
self._reader.start()
self._stderr_reader = threading.Thread(target=self._read_stderr, daemon=True)
self._reader.start()
self._stderr_reader.start()
# ---------- lifecycle ----------
def initialize(
self,
client_name: str = "hermes",
@@ -165,16 +115,11 @@ class CodexAppServerClient:
capabilities: Optional[dict] = None,
timeout: float = 10.0,
) -> dict:
"""Send `initialize` + `initialized` handshake. Returns the server's
InitializeResponse (userAgent, codexHome, platformFamily, platformOs)."""
"""Send ``initialize`` + ``initialized``; return the server's InitializeResponse."""
if self._initialized:
raise RuntimeError("already initialized")
params = {
"clientInfo": {
"name": client_name,
"title": client_title,
"version": client_version,
},
"clientInfo": {"name": client_name, "title": client_title, "version": client_version},
"capabilities": capabilities or {},
}
result = self.request("initialize", params, timeout=timeout)
@@ -208,16 +153,8 @@ class CodexAppServerClient:
def __exit__(self, *exc: Any) -> None:
self.close()
# ---------- send/receive ----------
def request(
self,
method: str,
params: Optional[dict] = None,
timeout: float = 30.0,
) -> dict:
"""Send a JSON-RPC request and block on the response. Returns `result`,
raises CodexAppServerError on `error`."""
def request(self, method: str, params: Optional[dict] = None, timeout: float = 30.0) -> dict:
"""Send a request and block for ``result``; raise CodexAppServerError on ``error``."""
rid = self._take_id()
q: queue.Queue = queue.Queue(maxsize=1)
with self._pending_lock:
@@ -228,16 +165,10 @@ class CodexAppServerClient:
except queue.Empty:
with self._pending_lock:
self._pending.pop(rid, None)
raise TimeoutError(
f"codex app-server method {method!r} timed out after {timeout}s"
)
raise TimeoutError(f"codex app-server method {method!r} timed out after {timeout}s")
if "error" in msg:
err = msg["error"]
raise CodexAppServerError(
code=err.get("code", -1),
message=err.get("message", ""),
data=err.get("data"),
)
raise CodexAppServerError(code=err.get("code", -1), message=err.get("message", ""), data=err.get("data"))
return msg.get("result", {})
def notify(self, method: str, params: Optional[dict] = None) -> None:
@@ -248,37 +179,29 @@ class CodexAppServerClient:
"""Reply to a server-initiated request (e.g. approval prompts)."""
self._send({"id": request_id, "result": result})
def respond_error(
self, request_id: Any, code: int, message: str, data: Optional[Any] = None
) -> None:
def respond_error(self, request_id: Any, code: int, message: str, data: Optional[Any] = None) -> None:
"""Reply to a server-initiated request with an error."""
err: dict[str, Any] = {"code": code, "message": message}
if data is not None:
err["data"] = data
self._send({"id": request_id, "error": err})
def take_notification(self, timeout: float = 0.0) -> Optional[dict]:
"""Pop the next streaming notification, or return None on timeout.
timeout=0.0 means non-blocking. Use small positive timeouts inside the
AIAgent turn loop to interleave reads with interrupt checks."""
@staticmethod
def _take(q: queue.Queue, timeout: float) -> Optional[dict]:
try:
if timeout <= 0:
return self._notifications.get_nowait()
return self._notifications.get(timeout=timeout)
return q.get_nowait()
return q.get(timeout=timeout)
except queue.Empty:
return None
def take_notification(self, timeout: float = 0.0) -> Optional[dict]:
"""Pop the next streaming notification, or None on timeout (0 = non-blocking)."""
return self._take(self._notifications, timeout)
def take_server_request(self, timeout: float = 0.0) -> Optional[dict]:
"""Pop the next server-initiated request (e.g. exec/applyPatch approval)."""
try:
if timeout <= 0:
return self._server_requests.get_nowait()
return self._server_requests.get(timeout=timeout)
except queue.Empty:
return None
# ---------- diagnostics ----------
return self._take(self._server_requests, timeout)
def stderr_tail(self, n: int = 20) -> list[str]:
"""Return last n lines of codex's stderr (for error reports)."""
@@ -288,12 +211,7 @@ class CodexAppServerClient:
def is_alive(self) -> bool:
return self._proc.poll() is None
# ---------- internals ----------
def _take_id(self) -> int:
# JSON-RPC ids only need to be unique per-connection. A simple
# monotonically increasing int is the common choice and matches what
# codex's own clients use.
rid = self._next_id
self._next_id += 1
return rid
@@ -307,38 +225,34 @@ class CodexAppServerClient:
self._proc.stdin.write((json.dumps(obj) + "\n").encode("utf-8"))
self._proc.stdin.flush()
except (BrokenPipeError, ValueError) as exc:
raise RuntimeError(
f"codex app-server stdin closed unexpectedly: {exc}"
) from exc
raise RuntimeError(f"codex app-server stdin closed unexpectedly: {exc}") from exc
def _append_stderr(self, line: str) -> None:
with self._stderr_lock:
self._stderr_lines.append(line)
if len(self._stderr_lines) > 500: # bound memory
self._stderr_lines = self._stderr_lines[-500:]
def _read_stdout(self) -> None:
if self._proc.stdout is None:
return
try:
for line in iter(self._proc.stdout.readline, b""):
if not line:
break
line = line.strip()
if not line:
continue
try:
msg = json.loads(line)
except json.JSONDecodeError:
# Non-JSON output is unexpected on stdout; tracing belongs
# on stderr. Surface it via stderr buffer for diagnostics.
with self._stderr_lock:
self._stderr_lines.append(
f"<non-json on stdout> {line[:200]!r}"
)
# Non-JSON stdout is unexpected; surface it via the stderr buffer.
self._append_stderr(f"<non-json on stdout> {line[:200]!r}")
continue
self._dispatch(msg)
except Exception as exc:
with self._stderr_lock:
self._stderr_lines.append(f"<stdout reader error> {exc}")
self._append_stderr(f"<stdout reader error> {exc}")
def _dispatch(self, msg: dict) -> None:
# Reply (has id + result/error, no method)
if "id" in msg and ("result" in msg or "error" in msg):
if "id" in msg and ("result" in msg or "error" in msg): # reply
with self._pending_lock:
pending = self._pending.pop(msg["id"], None)
if pending is not None:
@@ -346,63 +260,36 @@ class CodexAppServerClient:
pending.queue.put_nowait(msg)
except queue.Full: # pragma: no cover - defensive
pass
return
# Server-initiated request (has id + method)
if "id" in msg and "method" in msg:
self._server_requests.put(msg)
return
# Notification (no id)
if "method" in msg:
self._notifications.put(msg)
elif "method" in msg: # server-initiated request (has id) or notification
(self._server_requests if "id" in msg else self._notifications).put(msg)
def _read_stderr(self) -> None:
if self._proc.stderr is None:
return
try:
for line in iter(self._proc.stderr.readline, b""):
if not line:
break
with self._stderr_lock:
self._stderr_lines.append(
line.decode("utf-8", "replace").rstrip()
)
# Bound memory: keep last 500 lines.
if len(self._stderr_lines) > 500:
self._stderr_lines = self._stderr_lines[-500:]
self._append_stderr(line.decode("utf-8", "replace").rstrip())
except Exception: # pragma: no cover
pass
def parse_codex_version(output: str) -> Optional[tuple[int, int, int]]:
"""Parse `codex --version` output. Returns (major, minor, patch) or None."""
# Output format: "codex-cli 0.130.0" possibly followed by metadata.
import re
"""Parse ``codex --version`` output ("codex-cli 0.130.0 ...") into (major, minor, patch)."""
match = re.search(r"(\d+)\.(\d+)\.(\d+)", output or "")
if not match:
return None
return (int(match.group(1)), int(match.group(2)), int(match.group(3)))
return tuple(int(g) for g in match.groups()) if match else None
def check_codex_binary(
codex_bin: str = "codex", min_version: tuple[int, int, int] = MIN_CODEX_VERSION
) -> tuple[bool, str]:
"""Verify codex CLI is installed and meets minimum version.
Returns (ok, message). Used by setup wizard and runtime startup."""
"""Verify codex CLI is installed and meets minimum version. Returns (ok, message)."""
try:
proc = subprocess.run(
[codex_bin, "--version"],
capture_output=True,
text=True, encoding='utf-8', errors='replace',
timeout=10,
stdin=subprocess.DEVNULL,
[codex_bin, "--version"], capture_output=True, text=True, encoding='utf-8', errors='replace',
timeout=10, stdin=subprocess.DEVNULL,
)
except FileNotFoundError:
return False, (
f"codex CLI not found at {codex_bin!r}. Install with: "
f"npm i -g @openai/codex"
)
return False, f"codex CLI not found at {codex_bin!r}. Install with: npm i -g @openai/codex"
except subprocess.TimeoutExpired:
return False, "codex --version timed out"
if proc.returncode != 0:
@@ -410,9 +297,7 @@ def check_codex_binary(
version = parse_codex_version(proc.stdout)
if version is None:
return False, f"could not parse codex version from: {proc.stdout!r}"
have, need = ".".join(map(str, version)), ".".join(map(str, min_version))
if version < min_version:
return False, (
f"codex {'.'.join(map(str, version))} is older than required "
f"{'.'.join(map(str, min_version))}. Run: npm i -g @openai/codex"
)
return True, ".".join(map(str, version))
return False, f"codex {have} is older than required {need}. Run: npm i -g @openai/codex"
return True, have

File diff suppressed because it is too large Load Diff

View File

@@ -1,29 +1,14 @@
"""Projects codex app-server events into Hermes' messages list.
The translator that lets Hermes' memory/skill review keep working under the
Codex runtime: it converts Codex `item/*` notifications into the standard
OpenAI-shaped `{role, content, tool_calls, tool_call_id}` entries that
`agent/curator.py` already knows how to read.
Codex emits items with a discriminator field `type`:
- userMessage → {role: "user", content}
- agentMessage → {role: "assistant", content}
- reasoning → stashed in the assistant's "reasoning" field
- commandExecution → assistant tool_call(name="exec") + tool result
- fileChange → assistant tool_call(name="apply_patch") + tool result
- mcpToolCall → assistant tool_call(name=f"mcp.{server}.{tool}") + tool result
- dynamicToolCall → assistant tool_call(name=tool) + tool result
- plan/hookPrompt/collabAgentToolCall → recorded as opaque assistant notes
Each item maps to AT MOST one assistant entry + one tool entry, preserving
Hermes' message-alternation invariants (system → user → assistant → user/tool
→ assistant → ...). Multiple Codex tool calls within one Codex turn produce
multiple consecutive (assistant, tool) pairs, which is the same shape Hermes
already produces for parallel tool calls.
Counters tracked alongside projection:
- tool_iterations: ticks once per completed tool-shaped item. Used by
AIAgent._iters_since_skill (skill nudge gate, default threshold 10).
Converts Codex ``item/*`` notifications into OpenAI-shaped
``{role, content, tool_calls, tool_call_id}`` entries that memory/skill review
(agent/curator.py) already reads:
userMessage → user; agentMessage → assistant; reasoning → stashed onto the next
assistant entry; commandExecution / fileChange / mcpToolCall / dynamicToolCall →
assistant tool_call + tool result; anything else → opaque assistant note.
Each item yields AT MOST one assistant + one tool entry, preserving Hermes'
message-alternation invariants. ``is_tool_iteration`` ticks once per completed
tool-shaped item (AIAgent._iters_since_skill skill-nudge gate).
"""
from __future__ import annotations
@@ -31,20 +16,13 @@ from __future__ import annotations
import hashlib
import json
from dataclasses import dataclass, field
from typing import Any, Optional
from typing import Any, Callable, Optional
def _deterministic_call_id(item_type: str, item_id: str) -> str:
"""Stable id for tool_call message correlation.
Uses the codex item id directly when present (already a uuid); falls back
to a content hash so replay produces the same id across sessions and
prefix caches stay valid. See AGENTS.md Pitfall #16 (deterministic IDs in
tool call history)."""
if item_id:
return f"codex_{item_type}_{item_id}"
digest = hashlib.sha256(f"{item_type}".encode()).hexdigest()[:16]
return f"codex_{item_type}_{digest}"
"""Stable tool_call id: the codex item id when present, else a content hash so
replay yields the same id across sessions and prefix caches stay valid."""
return f"codex_{item_type}_{item_id or hashlib.sha256(f'{item_type}'.encode()).hexdigest()[:16]}"
def _format_tool_args(d: dict) -> str:
@@ -52,14 +30,18 @@ def _format_tool_args(d: dict) -> str:
return json.dumps(d, ensure_ascii=False, sort_keys=True)
def _dict_args(raw: Any) -> dict:
args = raw or {}
return args if isinstance(args, dict) else {"arguments": args}
@dataclass
class ProjectionResult:
"""Output of projecting one Codex item.
`messages` is a list because some Codex items produce two messages
(assistant tool_call + tool result). Empty list = item ignored (e.g. a
streaming `outputDelta` that doesn't materialize into messages until the
`item/completed` event)."""
``messages`` may hold two entries (assistant tool_call + tool result); empty
means the item was ignored (e.g. a streaming delta before ``item/completed``).
"""
messages: list[dict] = field(default_factory=list)
is_tool_iteration: bool = False
@@ -69,65 +51,56 @@ class ProjectionResult:
class CodexEventProjector:
"""Stateful projector consuming Codex notifications in arrival order.
Owns the in-progress reasoning content (codex emits reasoning as separate
items but Hermes stashes it on the next assistant message)."""
Owns in-progress reasoning: codex emits it as separate items, Hermes stashes
it on the next assistant message.
"""
def __init__(self) -> None:
self._pending_reasoning: list[str] = []
def project(self, notification: dict) -> ProjectionResult:
"""Project a single notification. Idempotent for non-completion events;
only `item/completed` and `turn/completed` materialize messages."""
"""Project one notification; only ``item/completed`` materializes messages.
Streaming deltas are display-only, mirroring how Hermes writes the
assistant message only after the streaming completion event.
"""
method = notification.get("method", "")
params = notification.get("params", {}) or {}
# We only materialize messages on `item/completed`. Streaming deltas
# (`item/<type>/outputDelta`, `item/<type>/delta`) are display-only and
# don't enter the messages list — same way Hermes already only writes
# the assistant message after the streaming completion event.
if method != "item/completed":
return ProjectionResult()
item = params.get("item") or {}
item_type = item.get("type") or ""
item_id = item.get("id") or ""
if item_type == "agentMessage":
return self._project_agent_message(item)
if item_type == "reasoning":
self._pending_reasoning.extend(item.get("summary") or [])
self._pending_reasoning.extend(item.get("content") or [])
return ProjectionResult()
if item_type == "commandExecution":
return self._project_command(item, item_id)
if item_type == "fileChange":
return self._project_file_change(item, item_id)
if item_type == "mcpToolCall":
return self._project_mcp_tool_call(item, item_id)
if item_type == "dynamicToolCall":
return self._project_dynamic_tool_call(item, item_id)
if item_type == "userMessage":
return self._project_user_message(item)
# Unknown / rare items (plan, hookPrompt, collabAgentToolCall, etc.)
# — record as opaque assistant note so memory review can still see
# *something* happened, but don't fabricate tool_call structure.
tool_projection = self._TOOL_PROJECTIONS.get(item_type)
if tool_projection is not None:
return self._project_tool_item(item, item_id, tool_projection)
# Unknown / rare items (plan, hookPrompt, collabAgentToolCall, ...): keep an
# opaque note so memory review sees *something* happened, without
# fabricating tool_call structure.
return self._project_opaque(item, item_type)
# ---------- per-type projections ----------
def _project_agent_message(self, item: dict) -> ProjectionResult:
text = item.get("text") or ""
msg: dict[str, Any] = {"role": "assistant", "content": text}
def _assistant_message(self, content: Optional[str], **extra: Any) -> dict[str, Any]:
msg: dict[str, Any] = {"role": "assistant", "content": content, **extra}
if self._pending_reasoning:
msg["reasoning"] = "\n".join(self._pending_reasoning)
self._pending_reasoning = []
return ProjectionResult(messages=[msg], final_text=text)
return msg
def _project_agent_message(self, item: dict) -> ProjectionResult:
text = item.get("text") or ""
return ProjectionResult(messages=[self._assistant_message(text)], final_text=text)
def _project_user_message(self, item: dict) -> ProjectionResult:
# codex's userMessage content is a list of UserInput variants. For
# projection purposes we flatten any text fragments and ignore
# non-text parts (images, etc.) — Hermes' messages store text only.
# userMessage content is a list of UserInput variants; flatten text
# fragments and drop non-text parts (Hermes' messages store text only).
text_parts: list[str] = []
for fragment in item.get("content") or []:
if isinstance(fragment, dict):
@@ -135,111 +108,54 @@ class CodexEventProjector:
text_parts.append(fragment.get("text") or "")
elif "text" in fragment:
text_parts.append(str(fragment["text"]))
return ProjectionResult(
messages=[{"role": "user", "content": "\n".join(text_parts)}]
)
return ProjectionResult(messages=[{"role": "user", "content": "\n".join(text_parts)}])
def _project_command(self, item: dict, item_id: str) -> ProjectionResult:
call_id = _deterministic_call_id("exec", item_id)
args = {
"command": item.get("command") or "",
"cwd": item.get("cwd") or "",
}
assistant_msg = {
"role": "assistant",
"content": None,
"tool_calls": [
{
"id": call_id,
"type": "function",
"function": {
"name": "exec_command",
"arguments": _format_tool_args(args),
},
}
],
}
if self._pending_reasoning:
assistant_msg["reasoning"] = "\n".join(self._pending_reasoning)
self._pending_reasoning = []
def _project_tool_item(
self,
item: dict,
item_id: str,
spec: Callable[[dict], tuple[str, str, dict, str]],
) -> ProjectionResult:
"""Emit the (assistant tool_call, tool result) pair for a tool-shaped item.
``spec(item)`` returns ``(call_id_type, tool_name, args, tool_content)``.
"""
id_type, name, args, content = spec(item)
call_id = _deterministic_call_id(id_type, item_id)
assistant_msg = self._assistant_message(
None,
tool_calls=[{"id": call_id, "type": "function", "function": {"name": name, "arguments": _format_tool_args(args)}}],
)
tool_msg = {"role": "tool", "tool_call_id": call_id, "content": content}
return ProjectionResult(messages=[assistant_msg, tool_msg], is_tool_iteration=True)
@staticmethod
def _command_spec(item: dict) -> tuple[str, str, dict, str]:
args = {"command": item.get("command") or "", "cwd": item.get("cwd") or ""}
output = item.get("aggregatedOutput") or ""
exit_code = item.get("exitCode")
if exit_code is not None and exit_code != 0:
output = f"[exit {exit_code}]\n{output}"
tool_msg = {
"role": "tool",
"tool_call_id": call_id,
"content": output,
}
return ProjectionResult(
messages=[assistant_msg, tool_msg], is_tool_iteration=True
)
return "exec", "exec_command", args, output
def _project_file_change(self, item: dict, item_id: str) -> ProjectionResult:
call_id = _deterministic_call_id("apply_patch", item_id)
# Reduce the codex changes array to a digest the agent loop will
# find readable. We record per-file change kinds (Add/Update/Delete)
# without inlining full file contents — those can be huge.
changes_summary = []
for change in item.get("changes") or []:
kind = (change.get("kind") or {}).get("type") or "update"
path = change.get("path") or ""
changes_summary.append({"kind": kind, "path": path})
args = {"changes": changes_summary}
assistant_msg = {
"role": "assistant",
"content": None,
"tool_calls": [
{
"id": call_id,
"type": "function",
"function": {
"name": "apply_patch",
"arguments": _format_tool_args(args),
},
}
],
}
if self._pending_reasoning:
assistant_msg["reasoning"] = "\n".join(self._pending_reasoning)
self._pending_reasoning = []
@staticmethod
def _file_change_spec(item: dict) -> tuple[str, str, dict, str]:
# Per-file change kinds only — full file contents can be huge.
changes_summary = [
{
"kind": (change.get("kind") or {}).get("type") or "update",
"path": change.get("path") or "",
}
for change in item.get("changes") or []
]
status = item.get("status") or "unknown"
n = len(changes_summary)
tool_msg = {
"role": "tool",
"tool_call_id": call_id,
"content": f"apply_patch status={status}, {n} change(s)",
}
return ProjectionResult(
messages=[assistant_msg, tool_msg], is_tool_iteration=True
)
content = f"apply_patch status={status}, {len(changes_summary)} change(s)"
return "apply_patch", "apply_patch", {"changes": changes_summary}, content
def _project_mcp_tool_call(self, item: dict, item_id: str) -> ProjectionResult:
@staticmethod
def _mcp_tool_call_spec(item: dict) -> tuple[str, str, dict, str]:
server = item.get("server") or "mcp"
tool = item.get("tool") or "unknown"
# Mirror the native MCP tool-name convention (mcp__server__tool) so the
# deterministic call_id input stays consistent with registration names.
call_id = _deterministic_call_id(f"mcp__{server}__{tool}", item_id)
args = item.get("arguments") or {}
if not isinstance(args, dict):
args = {"arguments": args}
assistant_msg = {
"role": "assistant",
"content": None,
"tool_calls": [
{
"id": call_id,
"type": "function",
"function": {
"name": f"mcp.{server}.{tool}",
"arguments": _format_tool_args(args),
},
}
],
}
if self._pending_reasoning:
assistant_msg["reasoning"] = "\n".join(self._pending_reasoning)
self._pending_reasoning = []
result = item.get("result")
error = item.get("error")
if error:
@@ -248,67 +164,30 @@ class CodexEventProjector:
content = json.dumps(result, ensure_ascii=False)[:4000]
else:
content = ""
tool_msg = {
"role": "tool",
"tool_call_id": call_id,
"content": content,
}
return ProjectionResult(
messages=[assistant_msg, tool_msg], is_tool_iteration=True
)
# Mirror the native MCP name convention (mcp__server__tool) in the call id
# so it stays consistent with registration names.
return f"mcp__{server}__{tool}", f"mcp.{server}.{tool}", _dict_args(item.get("arguments")), content
def _project_dynamic_tool_call(
self, item: dict, item_id: str
) -> ProjectionResult:
@staticmethod
def _dynamic_tool_call_spec(item: dict) -> tuple[str, str, dict, str]:
tool = item.get("tool") or "unknown"
call_id = _deterministic_call_id(f"dyn_{tool}", item_id)
args = item.get("arguments") or {}
if not isinstance(args, dict):
args = {"arguments": args}
assistant_msg = {
"role": "assistant",
"content": None,
"tool_calls": [
{
"id": call_id,
"type": "function",
"function": {
"name": tool,
"arguments": _format_tool_args(args),
},
}
],
}
if self._pending_reasoning:
assistant_msg["reasoning"] = "\n".join(self._pending_reasoning)
self._pending_reasoning = []
content_items = item.get("contentItems") or []
if isinstance(content_items, list) and content_items:
content = json.dumps(content_items, ensure_ascii=False)[:4000]
else:
success = item.get("success")
content = f"success={success}"
tool_msg = {
"role": "tool",
"tool_call_id": call_id,
"content": content,
}
return ProjectionResult(
messages=[assistant_msg, tool_msg], is_tool_iteration=True
)
content = f"success={item.get('success')}"
return f"dyn_{tool}", tool, _dict_args(item.get("arguments")), content
_TOOL_PROJECTIONS: dict[str, Callable[[dict], tuple[str, str, dict, str]]] = {
"commandExecution": _command_spec,
"fileChange": _file_change_spec,
"mcpToolCall": _mcp_tool_call_spec,
"dynamicToolCall": _dynamic_tool_call_spec,
}
def _project_opaque(self, item: dict, item_type: str) -> ProjectionResult:
# Record the existence of the item without inventing tool_calls.
# Memory review will see this and may or may not save anything.
try:
payload = json.dumps(item, ensure_ascii=False)[:1500]
except (TypeError, ValueError):
payload = repr(item)[:1500]
return ProjectionResult(
messages=[
{
"role": "assistant",
"content": f"[codex {item_type}] {payload}",
}
]
)
return ProjectionResult(messages=[{"role": "assistant", "content": f"[codex {item_type}] {payload}"}])

View File

@@ -1,45 +1,10 @@
"""Hermes-tools-as-MCP server for the codex_app_server runtime.
When the user runs `openai/*` turns through the codex app-server, codex
owns the loop and builds its own tool list. By default, that means
Hermes' richer tool surface — web search, browser automation,
delegate_task subagents, vision analysis, persistent memory, skills,
cross-session search, image generation, TTS — is unreachable.
This module exposes a curated subset of those Hermes tools to the
spawned codex subprocess via stdio MCP. Codex registers it as a normal
MCP server (per `~/.codex/config.toml [mcp_servers.hermes-tools]`) and
the user gets full Hermes capability inside a Codex turn.
Scope (what we expose):
- web_search, web_extract — Firecrawl, no codex equivalent
- browser_navigate / _click / _type / — Camofox/Browserbase automation
_snapshot / _scroll / _back / _press /
_get_images / _console / _vision
- vision_analyze — image inspection by vision model
- image_generate — image generation
- skill_view, skills_list — Hermes' skill library
- text_to_speech — TTS
- kanban_* (complete/block/comment/ — kanban worker + orchestrator
heartbeat/show/list/create/ handoff (stateless: read env var,
unblock/link) write ~/.hermes/kanban.db)
What we DO NOT expose:
- terminal / shell — codex's own shell tool
- read_file / write_file / patch — codex's apply_patch + shell
- search_files / process — codex's shell
- clarify — codex's own UX
- delegate_task / memory / — `_AGENT_LOOP_TOOLS` in Hermes
session_search / todo (model_tools.py). They require
the running AIAgent context to
dispatch (mid-loop state), so a
stateless MCP callback can't
drive them. See the inline
comment on EXPOSED_TOOLS below.
Run with: python -m agent.transports.hermes_tools_mcp_server
Spawned by: CodexAppServerSession.ensure_started() when the runtime is
active and config opts in.
Under the codex app-server, codex owns the loop and its own tool list, so
Hermes' richer surface (web search, browser, vision, image gen, skills, TTS,
kanban handoff) would be unreachable. This module exposes a curated subset over
stdio MCP; codex registers it via ``~/.codex/config.toml [mcp_servers.hermes-tools]``.
Run with: ``python -m agent.transports.hermes_tools_mcp_server``.
"""
from __future__ import annotations
@@ -54,61 +19,31 @@ from typing import Any, Optional
logger = logging.getLogger(__name__)
# JSON Schema type -> Python type mapping for signature generation
_JSON_TO_PY = {
"string": str,
"integer": int,
"number": float,
"boolean": bool,
"array": list,
"object": dict,
}
_JSON_TO_PY = {"string": str, "integer": int, "number": float, "boolean": bool, "array": list, "object": dict}
def _signature_from_schema(schema: dict | None) -> tuple[inspect.Signature, dict[str, type]]:
"""Build a Python function signature and annotations from a JSON schema.
Args:
schema: JSON Schema dict with "properties" and "required" keys.
Returns:
(signature, annotations_dict) where signature has KEYWORD_ONLY params
and annotations maps param names to Python types.
"""
"""Build a KEYWORD_ONLY signature + annotations dict from a JSON schema's
``properties`` / ``required`` (optional params default to None)."""
props = (schema or {}).get("properties") or {}
required = set((schema or {}).get("required") or [])
params, annots = [], {}
for pname, pspec in props.items():
if pname.startswith("_"):
continue
py = _JSON_TO_PY.get((pspec or {}).get("type"), Any)
ann, default = (
(py, inspect.Parameter.empty)
if pname in required
else (Optional[py], None)
)
ann, default = (py, inspect.Parameter.empty) if pname in required else (Optional[py], None)
annots[pname] = ann
params.append(
inspect.Parameter(
pname, inspect.Parameter.KEYWORD_ONLY, annotation=ann, default=default
)
)
params.append(inspect.Parameter(pname, inspect.Parameter.KEYWORD_ONLY, annotation=ann, default=default))
return inspect.Signature(params, return_annotation=str), annots
# Tools we expose. Each name MUST match a registered Hermes tool that
# `model_tools.handle_function_call()` can dispatch.
#
# What we deliberately DO NOT expose:
# - terminal / shell / read_file / write_file / patch / search_files /
# process — codex's built-ins cover these and approval routes through
# codex's own UI.
# - delegate_task / memory / session_search / todo — these are
# `_AGENT_LOOP_TOOLS` in Hermes (model_tools.py:493). They require
# the running AIAgent context to dispatch (mid-loop state), so a
# stateless MCP callback can't drive them. Hermes' default runtime
# keeps these working; the codex_app_server runtime cannot.
# Each name MUST match a registered Hermes tool that
# ``model_tools.handle_function_call()`` can dispatch.
# Deliberately NOT exposed: terminal/shell, read_file/write_file/patch,
# search_files/process, clarify — codex's built-ins cover them with codex's own
# approval UI; delegate_task/memory/session_search/todo — ``_AGENT_LOOP_TOOLS``
# need the running AIAgent context, which a stateless MCP callback lacks.
EXPOSED_TOOLS: tuple[str, ...] = (
"web_search",
"web_extract",
@@ -127,12 +62,8 @@ EXPOSED_TOOLS: tuple[str, ...] = (
"skill_view",
"skills_list",
"text_to_speech",
# Kanban worker handoff tools — gated on HERMES_KANBAN_TASK env var
# (set by the kanban dispatcher when spawning a worker). Without these
# in the callback, a worker spawned with openai_runtime=codex_app_server
# could do the work but couldn't report completion back to the kernel,
# making it hang until timeout. Stateless dispatch — they just read
# the env var and write to ~/.hermes/kanban.db.
# Kanban handoff tools: stateless (read HERMES_KANBAN_TASK, write kanban.db).
# Without them a codex-runtime worker can't report completion and hangs.
"kanban_complete",
"kanban_block",
"kanban_request_review",
@@ -141,10 +72,7 @@ EXPOSED_TOOLS: tuple[str, ...] = (
"kanban_heartbeat",
"kanban_show",
"kanban_list",
# NOTE: kanban_create / kanban_unblock / kanban_link are orchestrator-
# only — the kanban tool gates them on HERMES_KANBAN_TASK being unset.
# They're exposed here for orchestrator agents running on the codex
# runtime that need to dispatch new tasks.
# Orchestrator-only (the kanban tool gates them on HERMES_KANBAN_TASK unset).
"kanban_create",
"kanban_unblock",
"kanban_link",
@@ -152,23 +80,15 @@ EXPOSED_TOOLS: tuple[str, ...] = (
def _build_server() -> Any:
"""Create the MCP server with Hermes tools attached. Lazy imports
so the module can be imported without the mcp package installed
(we degrade to a clear error only when actually run)."""
"""Create the MCP server with Hermes tools attached (lazy imports so the module
imports without the mcp package; the clear error fires only when run)."""
try:
# mcp 2.0 removed `mcp.server.fastmcp`; `mcp.server.MCPServer` is the
# same decorator/add_tool surface under the new name.
# mcp 2.0 renamed `mcp.server.fastmcp` to `mcp.server.MCPServer` (same surface).
from mcp.server import MCPServer
except ImportError as exc: # pragma: no cover - install hint
raise ImportError(
f"hermes-tools MCP server requires the 'mcp' package: {exc}"
) from exc
raise ImportError(f"hermes-tools MCP server requires the 'mcp' package: {exc}") from exc
# Discover Hermes tools so dispatch works.
from model_tools import (
get_tool_definitions,
handle_function_call,
)
from model_tools import get_tool_definitions, handle_function_call
mcp = MCPServer(
"hermes-tools",
@@ -181,72 +101,49 @@ def _build_server() -> Any:
),
)
# Pull authoritative Hermes tool schemas for the ones we expose, so
# MCP clients see the same parameter docs Hermes gives the model.
# Authoritative Hermes schemas so MCP clients see the same parameter docs the model does.
all_defs = {
td["function"]["name"]: td["function"]
for td in (get_tool_definitions(quiet_mode=True) or [])
if isinstance(td, dict) and td.get("type") == "function"
}
exposed_count = 0
def _make_handler(tool_name: str, schema: dict | None, description: str):
# The SDK derives the input schema from the callable's signature (no
# inputSchema parameter), so synthesize __signature__ from the Hermes JSON Schema.
sig, annots = _signature_from_schema(schema)
def _dispatch(**kwargs: Any) -> str:
try:
# Drop None so unset optionals aren't forwarded to the handler.
return handle_function_call(tool_name, {k: v for k, v in kwargs.items() if v is not None})
except Exception as exc:
logger.exception("tool %s raised", tool_name)
return json.dumps({"error": str(exc), "tool": tool_name})
_dispatch.__name__ = tool_name
_dispatch.__doc__ = description
_dispatch.__signature__ = sig
_dispatch.__annotations__ = {**annots, "return": str}
return _dispatch
exposed_count = 0
for name in EXPOSED_TOOLS:
spec = all_defs.get(name)
if spec is None:
logger.debug(
"skipping %s — not registered in this Hermes process", name
)
logger.debug("skipping %s — not registered in this Hermes process", name)
continue
description = spec.get("description") or f"Hermes {name} tool"
params_schema = spec.get("parameters") or {"type": "object", "properties": {}}
# The SDK wants a Python callable and derives the input schema from
# its signature — there is no inputSchema parameter on either the
# decorator or add_tool(). So build a closure that takes the arguments
# dict, dispatches via handle_function_call, returns the result
# string, and carries a __signature__ synthesized from the Hermes
# JSON Schema (see _signature_from_schema) for the SDK to read.
def _make_handler(tool_name: str, schema: dict | None):
sig, annots = _signature_from_schema(schema)
def _dispatch(**kwargs: Any) -> str:
try:
# Filter out None values before dispatch so unset optionals
# aren't forwarded to the handler.
args = {k: v for k, v in kwargs.items() if v is not None}
return handle_function_call(tool_name, args or {})
except Exception as exc:
logger.exception("tool %s raised", tool_name)
return json.dumps({"error": str(exc), "tool": tool_name})
_dispatch.__name__ = tool_name
_dispatch.__doc__ = description
_dispatch.__signature__ = sig
_dispatch.__annotations__ = {**annots, "return": str}
return _dispatch
handler = _make_handler(name, params_schema, description)
try:
mcp.add_tool(
_make_handler(name, params_schema),
name=name,
description=description,
)
mcp.add_tool(handler, name=name, description=description)
except TypeError:
# Older mcp SDK signature — fall back to decorator-style. The
# synthesized __signature__ on the handler still drives schema
# generation there.
handler = _make_handler(name, params_schema)
handler = mcp.tool(name=name, description=description)(handler)
# Older mcp SDK: decorator-style registration; __signature__ still drives schema.
mcp.tool(name=name, description=description)(_make_handler(name, params_schema, description))
exposed_count += 1
logger.info(
"hermes-tools MCP server registered %d/%d tools",
exposed_count,
len(EXPOSED_TOOLS),
)
logger.info("hermes-tools MCP server registered %d/%d tools", exposed_count, len(EXPOSED_TOOLS))
return mcp
@@ -254,15 +151,12 @@ def main(argv: Optional[list[str]] = None) -> int:
"""Entry point for `python -m agent.transports.hermes_tools_mcp_server`."""
argv = argv or sys.argv[1:]
verbose = "--verbose" in argv or "-v" in argv
log_level = logging.INFO if verbose else logging.WARNING
logging.basicConfig(
level=log_level,
level=logging.INFO if verbose else logging.WARNING,
stream=sys.stderr, # MCP uses stdio for protocol — logs MUST go to stderr
format="%(asctime)s [%(levelname)s] %(name)s: %(message)s",
)
# Quiet mode: keep Hermes' own banners off stdout (which is the MCP wire).
# Keep Hermes' own banners off stdout (the MCP wire).
os.environ.setdefault("HERMES_QUIET", "1")
os.environ.setdefault("HERMES_REDACT_SECRETS", "true")
@@ -271,13 +165,10 @@ def main(argv: Optional[list[str]] = None) -> int:
except ImportError as exc:
sys.stderr.write(f"hermes-tools MCP server cannot start: {exc}\n")
return 2
# MCPServer.run() defaults to stdio transport, which is what codex
# spawns us on.
try:
server.run()
server.run() # defaults to stdio transport, which codex spawns us on
except KeyboardInterrupt:
return 0
pass
except Exception as exc:
logger.exception("hermes-tools MCP server crashed")
sys.stderr.write(f"hermes-tools MCP server error: {exc}\n")

View File

@@ -1,11 +1,8 @@
"""Shared types for normalized provider responses.
These dataclasses define the canonical shape that all provider adapters
normalize responses to. The shared surface is intentionally minimal —
only fields that every downstream consumer reads are top-level.
Protocol-specific state goes in ``provider_data`` dicts (response-level
and per-tool-call) so that protocol-aware code paths can access it
without polluting the shared type.
Only fields every downstream consumer reads are top-level; protocol-specific
state lives in ``provider_data`` (response-level and per-tool-call) so
protocol-aware code can reach it without widening the shared type.
"""
from __future__ import annotations
@@ -19,17 +16,11 @@ from typing import Any
class ToolCall:
"""A normalized tool call from any provider.
``id`` is the protocol's canonical identifier — what gets used in
``tool_call_id`` / ``tool_use_id`` when constructing tool result
messages. May be ``None`` when the provider omits it; the agent
fills it via ``_deterministic_call_id()`` before storing in history.
``provider_data`` carries per-tool-call protocol metadata that only
protocol-aware code reads:
* Codex: ``{"call_id": "call_XXX", "response_item_id": "fc_XXX"}``
* Gemini: ``{"extra_content": {"google": {"thought_signature": "..."}}}``
* Others: ``None``
``id`` is the protocol's canonical identifier (``tool_call_id`` / ``tool_use_id``);
may be ``None`` when the provider omits it — the agent fills it via
``_deterministic_call_id()`` before storing history.
``provider_data``: Codex ``{"call_id", "response_item_id"}``, Gemini
``{"extra_content": {"google": {"thought_signature": ...}}}``, else ``None``.
"""
id: str | None
@@ -37,43 +28,23 @@ class ToolCall:
arguments: str # JSON string
provider_data: dict[str, Any] | None = field(default=None, repr=False)
# ── Backward compatibility ──────────────────────────────────
# The agent loop reads tc.function.name / tc.function.arguments
# throughout run_agent.py (45+ sites). These properties let
# NormalizedResponse pass through without the _nr_to_assistant_message
# shim, while keeping ToolCall's canonical fields flat.
# Back-compat: run_agent reads tc.function.name / tc.function.arguments (45+
# sites) and getattr()s the provider fields, so expose them as properties.
@property
def type(self) -> str:
return "function"
@property
def function(self) -> ToolCall:
"""Return self so tc.function.name / tc.function.arguments work."""
return self
@property
def call_id(self) -> str | None:
"""Codex call_id from provider_data, accessed via getattr by _build_assistant_message."""
return (self.provider_data or {}).get("call_id")
def _pd(self, key: str) -> Any:
return (self.provider_data or {}).get(key)
@property
def response_item_id(self) -> str | None:
"""Codex response_item_id from provider_data."""
return (self.provider_data or {}).get("response_item_id")
@property
def extra_content(self) -> dict[str, Any] | None:
"""Gemini extra_content (thought_signature) from provider_data.
Gemini 3 thinking models attach ``extra_content`` with a
``thought_signature`` to each tool call. This signature must be
replayed on subsequent API calls — without it the API rejects the
request with HTTP 400. The chat_completions transport stores this
in ``provider_data["extra_content"]``; this property exposes it so
``_build_assistant_message`` can ``getattr(tc, "extra_content")``
uniformly.
"""
return (self.provider_data or {}).get("extra_content")
call_id = property(lambda self: self._pd("call_id"))
response_item_id = property(lambda self: self._pd("response_item_id"))
# Gemini thought_signature; must be replayed on later calls or the API returns HTTP 400.
extra_content = property(lambda self: self._pd("extra_content"))
@dataclass
@@ -85,20 +56,22 @@ class Usage:
total_tokens: int = 0
cached_tokens: int = 0
@classmethod
def from_openai(cls, u: Any) -> Usage:
"""Build from an OpenAI-shaped usage object, treating missing/None counts as 0."""
return cls(
prompt_tokens=getattr(u, "prompt_tokens", 0) or 0,
completion_tokens=getattr(u, "completion_tokens", 0) or 0,
total_tokens=getattr(u, "total_tokens", 0) or 0,
)
@dataclass
class NormalizedResponse:
"""Normalized API response from any provider.
Shared fields are truly cross-provider — every caller can rely on
them without branching on api_mode. Protocol-specific state goes in
``provider_data`` so that only protocol-aware code paths read it.
Response-level ``provider_data`` examples:
* Anthropic: ``{"reasoning_details": [...]}``
* Codex: ``{"codex_reasoning_items": [...], "codex_message_items": [...]}``
* Others: ``None``
Response-level ``provider_data``: Anthropic ``{"reasoning_details": [...]}``,
Codex ``{"codex_reasoning_items": [...], "codex_message_items": [...]}``, else ``None``.
"""
content: str | None
@@ -108,73 +81,27 @@ class NormalizedResponse:
usage: Usage | None = None
provider_data: dict[str, Any] | None = field(default=None, repr=False)
# ── Backward compatibility ──────────────────────────────────
# The shim _nr_to_assistant_message() mapped these from provider_data.
# These properties let NormalizedResponse pass through directly.
@property
def reasoning_content(self) -> str | None:
pd = self.provider_data or {}
return pd.get("reasoning_content")
# Back-compat accessors so NormalizedResponse passes through where the old
# _nr_to_assistant_message() shim mapped these from provider_data.
def _pd(self, key: str) -> Any:
return (self.provider_data or {}).get(key)
@property
def reasoning_details(self):
pd = self.provider_data or {}
return pd.get("reasoning_details")
@property
def anthropic_content_blocks(self):
"""Verbatim, order-preserving Anthropic content blocks for a turn.
Present only when an Anthropic turn interleaves signed thinking with
tool_use — the one shape the parallel reasoning_details + tool_calls
lists reconstruct in the wrong order, invalidating thinking-block
signatures on replay. See agent/transports/anthropic.py.
"""
pd = self.provider_data or {}
return pd.get("anthropic_content_blocks")
@property
def bedrock_content_blocks(self):
"""Verbatim, order-preserving Bedrock Converse content blocks."""
pd = self.provider_data or {}
return pd.get("bedrock_content_blocks")
@property
def codex_reasoning_items(self):
pd = self.provider_data or {}
return pd.get("codex_reasoning_items")
@property
def codex_message_items(self):
pd = self.provider_data or {}
return pd.get("codex_message_items")
reasoning_content = property(lambda self: self._pd("reasoning_content"))
reasoning_details = property(lambda self: self._pd("reasoning_details"))
# Order-preserving Anthropic blocks, present only when a turn interleaves signed
# thinking with tool_use (replay order invalidates signatures otherwise).
anthropic_content_blocks = property(lambda self: self._pd("anthropic_content_blocks"))
bedrock_content_blocks = property(lambda self: self._pd("bedrock_content_blocks")) # order-preserving Converse blocks
codex_reasoning_items = property(lambda self: self._pd("codex_reasoning_items"))
codex_message_items = property(lambda self: self._pd("codex_message_items"))
# ---------------------------------------------------------------------------
# Factory helpers
# ---------------------------------------------------------------------------
def build_tool_call(
id: str | None,
name: str,
arguments: Any,
**provider_fields: Any,
) -> ToolCall:
"""Build a ``ToolCall``, auto-serialising *arguments* if it's a dict.
Any extra keyword arguments are collected into ``provider_data``.
"""
def build_tool_call(id: str | None, name: str, arguments: Any, **provider_fields: Any) -> ToolCall:
"""Build a ``ToolCall``; dict *arguments* are JSON-serialised, extra kwargs become ``provider_data``."""
args_str = json.dumps(arguments) if isinstance(arguments, dict) else str(arguments)
pd = dict(provider_fields) if provider_fields else None
return ToolCall(id=id, name=name, arguments=args_str, provider_data=pd)
return ToolCall(id=id, name=name, arguments=args_str, provider_data=dict(provider_fields) if provider_fields else None)
def map_finish_reason(reason: str | None, mapping: dict[str, str]) -> str:
"""Translate a provider-specific stop reason to the normalised set.
Falls back to ``"stop"`` for unknown or ``None`` reasons.
"""
if reason is None:
return "stop"
return mapping.get(reason, "stop")
"""Translate a provider stop reason via *mapping*; unknown or ``None`` -> ``"stop"``."""
return "stop" if reason is None else mapping.get(reason, "stop")