refactor(agent/anthropic_message_convert): share _strip_thinking, extract latest-turn thinking filter, compact docstrings; endpoints docstrings

This commit is contained in:
Teknium
2026-09-02 19:03:50 -07:00
parent eba201ff6e
commit 5eb0f4f7d1
2 changed files with 188 additions and 298 deletions
+39 -71
View File
@@ -1,14 +1,11 @@
"""Endpoint-family detection for Anthropic-compatible base URLs.
A dozen services speak the Anthropic Messages API but differ in auth style,
accepted beta headers, and request quirks (MiniMax, Kimi/Moonshot, DeepSeek,
OpenCode, Azure AI Foundry, Nous Portal, Bedrock). Every such difference is
decided from the configured base URL, so the predicates live together here.
Pure functions over a base-URL string - no I/O, no SDK, no credentials - so
both ``agent/anthropic_adapter.py`` and ``agent/anthropic_message_convert.py``
can depend on this module without a cycle. ``agent.anthropic_adapter``
re-exports every name below.
A dozen services speak the Anthropic Messages API but differ in auth style, accepted beta
headers, and request quirks (MiniMax, Kimi/Moonshot, DeepSeek, OpenCode, Azure AI Foundry, Nous
Portal, Bedrock). Every such difference is decided from the configured base URL, so the
predicates live together here as pure functions (no I/O, SDK or credentials) that both
``agent/anthropic_adapter.py`` and ``agent/anthropic_message_convert.py`` can import without a
cycle. ``agent.anthropic_adapter`` re-exports every name below.
"""
from urllib.parse import urlparse
@@ -20,9 +17,7 @@ _MINIMAX_ANTHROPIC_PREFIXES = ("https://api.minimax.io/anthropic", "https://api.
def _normalize_base_url_text(base_url) -> str:
"""Coerce a base URL (str or ``httpx.URL``) to a stripped string; "" when falsy."""
if not base_url:
return ""
return str(base_url).strip()
return str(base_url).strip() if base_url else ""
def _normalized_lower(base_url) -> str:
@@ -31,39 +26,30 @@ def _normalized_lower(base_url) -> str:
def _is_third_party_anthropic_endpoint(base_url: str | None) -> bool:
"""True for any non-anthropic.com endpoint (own API keys via x-api-key; skip OAuth detection).
No base_url means the direct Anthropic API.
"""
"""Any non-anthropic.com endpoint (own x-api-key keys; skip OAuth detection). No base_url =
direct Anthropic API."""
normalized = _normalized_lower(base_url)
return bool(normalized) and "anthropic.com" not in normalized
def _is_kimi_coding_endpoint(base_url: str | None) -> bool:
"""True for Kimi's /coding endpoint, which requires a claude-code User-Agent."""
"""Kimi's /coding endpoint, which requires a claude-code User-Agent."""
return _normalized_lower(base_url).startswith("https://api.kimi.com/coding")
def _is_opencode_endpoint(base_url: str | None) -> bool:
"""True for OpenCode's Zen/Go relay (opencode.ai)."""
"""OpenCode's Zen/Go relay (opencode.ai)."""
return base_url_host_matches(base_url or "", "opencode.ai")
# Model-name prefixes identifying the Kimi / Moonshot family: official slugs
# (``kimi-k2.5``, ``kimi_thinking``, ``moonshot-v1-8k``) and release lines
# (``k1.5-…``, ``k2-thinking``, ``k25-…``, ``k3.x``/``k3-…``). Matched
# case-insensitively after stripping any ``vendor/`` prefix.
# Kimi / Moonshot family model-name prefixes: official slugs (``kimi-k2.5``, ``kimi_thinking``,
# ``moonshot-v1-8k``) and release lines (``k1.5-…``, ``k2-thinking``, ``k25-…``, ``k3.x``/``k3-…``).
# Matched case-insensitively after stripping any ``vendor/`` prefix.
_KIMI_FAMILY_MODEL_PREFIXES = (
"kimi-", "kimi_",
"moonshot-", "moonshot_",
"k1.", "k1-",
"k2.", "k2-",
"k25", "k2.5",
"k3.", "k3-",
"kimi-", "kimi_", "moonshot-", "moonshot_", "k1.", "k1-", "k2.", "k2-", "k25", "k2.5", "k3.", "k3-",
)
# Bare release slugs with no separator suffix (Kimi Coding Plan serves K3 as
# exactly ``k3``). Exact-match so unrelated names sharing the prefix don't match.
# Bare release slugs with no separator suffix (Kimi Coding Plan serves K3 as exactly ``k3``);
# exact-match so unrelated names sharing the prefix don't match.
_KIMI_FAMILY_EXACT_SLUGS = frozenset({"k3"})
@@ -79,14 +65,10 @@ def _model_name_is_kimi_family(model: str | None) -> bool:
def _is_kimi_family_endpoint(base_url: str | None, model: str | None = None) -> bool:
"""True for any Kimi / Moonshot Anthropic-Messages endpoint.
Broader than ``_is_kimi_coding_endpoint``: also matches any api.kimi.com /
moonshot.ai / moonshot.cn host, and any endpoint (e.g. a private gateway)
whose *model* is in the Kimi family — the upstream still enforces Kimi's
thinking semantics regardless of hostname. Decides whether unsigned
reasoning_content-derived thinking blocks are preserved on replay.
"""
"""Any Kimi / Moonshot Anthropic-Messages endpoint: the /coding endpoint, any api.kimi.com /
moonshot.ai / moonshot.cn host, or any endpoint (e.g. a private gateway) whose *model* is in
the Kimi family — the upstream enforces Kimi's thinking semantics regardless of hostname.
Decides whether unsigned reasoning_content-derived thinking blocks are preserved on replay."""
if _is_kimi_coding_endpoint(base_url):
return True
if any(base_url_host_matches(base_url or "", d) for d in ("api.kimi.com", "moonshot.ai", "moonshot.cn")):
@@ -95,27 +77,20 @@ def _is_kimi_family_endpoint(base_url: str | None, model: str | None = None) ->
def _is_deepseek_anthropic_endpoint(base_url: str | None) -> bool:
"""True for DeepSeek's ``/anthropic`` route.
In thinking mode DeepSeek requires prior-turn ``thinking`` blocks to round-trip
("The content[].thinking in the thinking mode must be passed back to the API"),
while the generic third-party path strips them. Its blocks are unsigned, so it
gets the same strip-signed / keep-unsigned policy as Kimi. Pinned to the
``/anthropic`` path so the OpenAI-compatible base URL is not misclassified.
"""
"""DeepSeek's ``/anthropic`` route. In thinking mode DeepSeek requires prior-turn ``thinking``
blocks to round-trip while the generic third-party path strips them; its blocks are unsigned,
so it gets the same strip-signed / keep-unsigned policy as Kimi. Pinned to the ``/anthropic``
path so the OpenAI-compatible base URL is not misclassified."""
if not base_url_host_matches(base_url or "", "api.deepseek.com"):
return False
return "/anthropic" in _normalized_lower(base_url)
def _is_nous_portal_endpoint(base_url: str | None) -> bool:
"""True for Nous Portal's Anthropic Messages route (Bearer JWT, verbatim catalog
ids, native thinking-signature replay).
Trusted hosts only: prod ``inference-api.nousresearch.com`` or the operator-set
``NOUS_INFERENCE_BASE_URL`` host (exact hostname equality, so neither lookalike
domains nor sibling hosts of the override match).
"""
"""Nous Portal's Anthropic Messages route (Bearer JWT, verbatim catalog ids, native
thinking-signature replay). Trusted hosts only: prod ``inference-api.nousresearch.com`` or the
operator-set ``NOUS_INFERENCE_BASE_URL`` host (exact hostname equality, so neither lookalike
domains nor sibling hosts of the override match)."""
if base_url_host_matches(base_url or "", "inference-api.nousresearch.com"):
return True
try:
@@ -131,13 +106,9 @@ def _is_nous_portal_endpoint(base_url: str | None) -> bool:
def _requires_bearer_auth(base_url: str | None) -> bool:
"""True for Anthropic-compatible providers that need ``Authorization: Bearer``
instead of ``x-api-key``: MiniMax, Azure AI Foundry, Palantir Foundry's LLM
proxy, CommandCode, and Nous Portal.
Palantir/CommandCode use hostname matching (not substring) so e.g.
``evil.com/palantirfoundry`` paths don't trigger Bearer auth.
"""
"""Providers needing ``Authorization: Bearer`` instead of ``x-api-key``: MiniMax, Azure AI
Foundry, Palantir Foundry's LLM proxy, CommandCode, Nous Portal. Palantir/CommandCode use
hostname matching (not substring) so ``evil.com/palantirfoundry`` paths don't trigger it."""
if _is_nous_portal_endpoint(base_url):
return True
normalized = _normalized_lower(base_url)
@@ -152,24 +123,21 @@ def _requires_bearer_auth(base_url: str | None) -> bool:
def _base_url_needs_context_1m_beta(base_url: str | None) -> bool:
"""True for endpoints that still gate 1M context behind a beta (Azure)."""
"""Endpoints that still gate 1M context behind a beta (Azure)."""
return "azure.com" in _normalize_base_url_text(base_url).lower()
def _is_minimax_anthropic_endpoint(base_url: str | None) -> bool:
"""True for MiniMax's Anthropic-compatible endpoints, which reject the
fine-grained-tool-streaming and context-1m betas (stripped even though MiniMax
also uses Bearer auth)."""
"""MiniMax's Anthropic-compatible endpoints, which reject the fine-grained-tool-streaming and
context-1m betas (stripped even though MiniMax also uses Bearer auth)."""
return _normalized_lower(base_url).startswith(_MINIMAX_ANTHROPIC_PREFIXES)
def _is_azure_anthropic_endpoint(base_url: str | None) -> bool:
"""True for Azure-hosted Anthropic Messages endpoints serving ``/anthropic``:
modern Foundry (``*.services.ai.azure.*``) and legacy Azure OpenAI
(``*.openai.azure.*``) hosts. Opts them into ``api-version`` query plumbing.
Deliberately no finite TLD allow-list, so sovereign/private clouds work.
"""
"""Azure-hosted Anthropic Messages endpoints serving ``/anthropic``: modern Foundry
(``*.services.ai.azure.*``) and legacy Azure OpenAI (``*.openai.azure.*``) hosts; opts them
into ``api-version`` query plumbing. Deliberately no finite TLD allow-list, so
sovereign/private clouds work."""
normalized = _normalize_base_url_text(base_url)
if not normalized:
return False
+149 -227
View File
@@ -1,14 +1,10 @@
"""OpenAI-style -> Anthropic Messages API request conversion.
Everything here rewrites *request payloads*: model-id normalization, tool
schemas, and the message list (content blocks, thinking blocks and their
signatures, tool_use/tool_result pairing, cache_control placement, screenshot
eviction, blank-block scrubbing).
Split out of ``agent/anthropic_adapter.py`` so the adapter keeps client
construction and the API call itself. Endpoint predicates come from
``agent/anthropic_endpoints.py``, so this module never imports the adapter and
there is no cycle. ``agent.anthropic_adapter`` re-exports every name below.
Everything here rewrites *request payloads*: model-id normalization, tool schemas, and the
message list (content blocks, thinking blocks and their signatures, tool_use/tool_result
pairing, cache_control placement, screenshot eviction, blank-block scrubbing). Endpoint
predicates come from ``agent/anthropic_endpoints.py``, so this module never imports the adapter
and there is no cycle. ``agent.anthropic_adapter`` re-exports every name below.
"""
import copy
@@ -18,9 +14,7 @@ import re
from typing import Any, Dict, List, Optional, Tuple
from agent.anthropic_endpoints import (
_is_deepseek_anthropic_endpoint,
_is_kimi_family_endpoint,
_is_nous_portal_endpoint,
_is_deepseek_anthropic_endpoint, _is_kimi_family_endpoint, _is_nous_portal_endpoint,
_is_third_party_anthropic_endpoint,
)
@@ -29,14 +23,13 @@ logger = logging.getLogger(__name__)
_THINKING_TYPES = frozenset(("thinking", "redacted_thinking"))
_CACHEABLE_TYPES = frozenset(("text", "tool_use"))
_EMPTY_TEXT_PLACEHOLDER = "(empty)"
_EMPTY_SCHEMA = {"type": "object", "properties": {}}
_BEDROCK_REGION_PREFIXES = (
"global.", "us.", "eu.", "apac.", "ap.", "au.", "jp.", "ca.", "sa.", "me.", "af.",
)
# ---------------------------------------------------------------------------
# Small shared predicates
# ---------------------------------------------------------------------------
# ----- small shared predicates -----
def _block_type(b: Any) -> Any:
@@ -49,12 +42,8 @@ def _has_block_type(blocks: List[Any], types) -> bool:
def _is_blank_text_block(b: Any) -> bool:
"""A text block whose ``text`` is not a non-whitespace string (None/int/blank all count).
Anthropic rejects such blocks with HTTP 400 ("text content blocks must contain
non-whitespace text"); checking isinstance first keeps a non-string from
reaching ``.strip()``.
"""
"""A text block whose ``text`` is not a non-whitespace string (None/int/blank all count) —
Anthropic 400s on them ("text content blocks must contain non-whitespace text")."""
if _block_type(b) != "text":
return False
text = b.get("text")
@@ -78,26 +67,24 @@ def _parse_tool_args(raw: Any) -> Any:
return {}
# ---------------------------------------------------------------------------
# Model / tool conversion
# ---------------------------------------------------------------------------
def _strip_thinking(blocks: List[Any]) -> List[Any]:
return [b for b in blocks if _block_type(b) not in _THINKING_TYPES]
# ----- model / tool conversion -----
def _is_bedrock_model_id(model: str) -> bool:
"""Bedrock ids (``anthropic.claude-opus-4-7``, ``us.anthropic.claude-*``) use dots
as namespace separators that must be preserved verbatim."""
"""Bedrock ids (``anthropic.claude-opus-4-7``, ``us.anthropic.claude-*``) use dots as namespace
separators that must be preserved verbatim."""
return model.lower().startswith(_BEDROCK_REGION_PREFIXES + ("anthropic.",))
def normalize_model_name(model: str, preserve_dots: bool = False) -> str:
"""Normalize a model name for the Anthropic API.
Strips the ``anthropic/`` prefix (case-insensitive) and, unless
``preserve_dots`` (DashScope: ``qwen3.5-plus``), converts version dots to
hyphens for Claude models only (``claude-opus-4.6`` -> ``claude-opus-4-6``).
Bedrock ids keep their namespace dots; non-Anthropic names (``gpt-5.4``)
keep dots as part of their canonical form.
"""
"""Strip the ``anthropic/`` prefix (case-insensitive) and, unless ``preserve_dots`` (DashScope:
``qwen3.5-plus``), convert version dots to hyphens for Claude models only (``claude-opus-4.6``
-> ``claude-opus-4-6``). Bedrock ids keep their namespace dots; non-Anthropic names
(``gpt-5.4``) keep dots as part of their canonical form."""
if model.lower().startswith("anthropic/"):
model = model[len("anthropic/"):]
if not preserve_dots and not _is_bedrock_model_id(model) and model.lower().startswith(("claude-", "anthropic/")):
@@ -113,22 +100,19 @@ def _sanitize_tool_id(tool_id: str) -> str:
def _normalize_tool_input_schema(schema: Any) -> Dict[str, Any]:
"""Normalize a tool schema for Anthropic's validator.
Collapses nullable unions (``anyOf: [{type: string}, {type: null}]``, which
Pydantic/MCP emit for optional fields) to the non-null branch — optionality is
already expressed by ``required``. ``keep_nullable_hint=False`` because the
OpenAPI ``nullable`` keyword is not recognized. Top-level oneOf/allOf/anyOf are
rejected with a generic 400, so they are dropped in favour of a plain object.
"""
"""Normalize a tool schema for Anthropic's validator: collapse nullable unions (``anyOf:
[{type: string}, {type: null}]`` from Pydantic/MCP optional fields) to the non-null branch —
optionality is already expressed by ``required``; ``keep_nullable_hint=False`` because the
OpenAPI ``nullable`` keyword is not recognized. Top-level oneOf/allOf/anyOf are rejected with a
generic 400, so they are dropped in favour of a plain object."""
if not schema:
return {"type": "object", "properties": {}}
return dict(_EMPTY_SCHEMA)
from tools.schema_sanitizer import strip_nullable_unions
normalized = strip_nullable_unions(schema, keep_nullable_hint=False)
if not isinstance(normalized, dict):
return {"type": "object", "properties": {}}
return dict(_EMPTY_SCHEMA)
banned = {"oneOf", "allOf", "anyOf"}
if banned & normalized.keys():
normalized = {k: v for k, v in normalized.items() if k not in banned}
@@ -139,11 +123,8 @@ def _normalize_tool_input_schema(schema: Any) -> Dict[str, Any]:
def convert_tools_to_anthropic(tools: List[Dict]) -> List[Dict]:
"""Convert OpenAI tool definitions to Anthropic format.
Duplicate names are dropped with a warning (Anthropic hard-400s on them);
``cache_control`` on the OpenAI tool dict is forwarded.
"""
"""Convert OpenAI tool definitions to Anthropic format. Duplicate names are dropped with a
warning (Anthropic hard-400s on them); ``cache_control`` on the OpenAI tool dict is forwarded."""
if not tools:
return []
result = []
@@ -152,9 +133,7 @@ def convert_tools_to_anthropic(tools: List[Dict]) -> List[Dict]:
fn = t.get("function", {})
name = fn.get("name", "")
if name and name in seen_names:
logger.warning(
"convert_tools_to_anthropic: duplicate tool name '%s' — dropping second occurrence", name
)
logger.warning("convert_tools_to_anthropic: duplicate tool name '%s' — dropping second occurrence", name)
continue
if name:
seen_names.add(name)
@@ -170,9 +149,7 @@ def convert_tools_to_anthropic(tools: List[Dict]) -> List[Dict]:
return result
# ---------------------------------------------------------------------------
# Content-part conversion
# ---------------------------------------------------------------------------
# ----- content-part conversion -----
def _image_source_from_openai_url(url: str) -> Dict[str, str]:
@@ -197,9 +174,8 @@ def _convert_content_part_to_anthropic(part: Any) -> Optional[Dict[str, Any]]:
ptype = part.get("type")
if ptype in ("input_text", "text"):
# Rebuild from whitelisted fields only: stored SDK text blocks carry
# output-only siblings (parsed_output, citations=None) that the INPUT
# schema rejects with 400 "Extra inputs are not permitted".
# Rebuild from whitelisted fields only: stored SDK text blocks carry output-only siblings
# (parsed_output, citations=None) that the INPUT schema rejects with 400.
block: Dict[str, Any] = _text_block(part.get("text", ""))
cits = part.get("citations")
if ptype == "text" and isinstance(cits, list) and cits:
@@ -218,12 +194,9 @@ def _convert_content_part_to_anthropic(part: Any) -> Optional[Dict[str, Any]]:
def _to_plain_data(value: Any, *, _depth: int = 0, _path: Optional[set] = None) -> Any:
"""Recursively convert SDK objects to plain Python data.
``_path`` tracks ids on the *current* recursion path (so shared but
non-cyclic objects convert normally while true cycles stringify); depth is
capped at 20.
"""
"""Recursively convert SDK objects to plain Python data. ``_path`` tracks ids on the *current*
recursion path (shared but non-cyclic objects convert normally while true cycles stringify);
depth is capped at 20."""
if _depth > 20:
return str(value)
if _path is None:
@@ -238,8 +211,8 @@ def _to_plain_data(value: Any, *, _depth: int = 0, _path: Optional[set] = None)
_path.add(obj_id)
if hasattr(value, "model_dump"):
try:
# warnings=False: streaming-accumulator blocks trip pydantic's
# serializer-mismatch UserWarning, which otherwise leaks to the terminal.
# warnings=False: streaming-accumulator blocks trip pydantic's serializer-mismatch
# UserWarning, which otherwise leaks to the terminal.
dumped = value.model_dump(warnings=False)
except TypeError: # duck-typed model_dump without pydantic's signature
dumped = value.model_dump()
@@ -276,8 +249,8 @@ def _convert_content_to_anthropic(content: Any) -> Any:
def _content_parts_to_anthropic_blocks(parts: Any) -> List[Dict[str, Any]]:
"""Tool-message content parts -> tool_result inner blocks (text + image only,
the types Anthropic accepts there). Used for multimodal tool results."""
"""Tool-message content parts -> tool_result inner blocks (text + image only, the types
Anthropic accepts there). Used for multimodal tool results."""
if not isinstance(parts, list):
return []
out: List[Dict[str, Any]] = []
@@ -293,12 +266,9 @@ def _content_parts_to_anthropic_blocks(parts: Any) -> List[Dict[str, Any]]:
def _safe_text(text: Any) -> str:
"""``text`` if non-whitespace, else the placeholder.
A blank text block stored in history (e.g. by compression) is replayed on every
turn and wedges the session with HTTP 400; the placeholder is self-healing.
Mirrors ``bedrock_adapter._safe_text`` (kept separate on purpose).
"""
"""``text`` if non-whitespace, else the placeholder. A blank text block stored in history (e.g.
by compression) is replayed on every turn and wedges the session with HTTP 400; the placeholder
is self-healing. Mirrors ``bedrock_adapter._safe_text`` (kept separate on purpose)."""
if text is None:
return _EMPTY_TEXT_PLACEHOLDER
if not isinstance(text, str):
@@ -306,15 +276,13 @@ def _safe_text(text: Any) -> str:
return text if text.strip() else _EMPTY_TEXT_PLACEHOLDER
# ---------------------------------------------------------------------------
# Replay-block sanitizing (per-type whitelist)
# ---------------------------------------------------------------------------
# ----- replay-block sanitizing (per-type whitelist) -----
def _replay_text(b: Dict[str, Any]) -> Optional[Dict[str, Any]]:
# Drop blank blocks rather than coerce in place: the caller relocates any
# cache_control and falls back to a placeholder only when nothing survives,
# so "(empty)" never sits as model-visible noise next to real blocks.
# Drop blank blocks rather than coerce in place: the caller relocates any cache_control and
# falls back to a placeholder only when nothing survives, so "(empty)" never sits as
# model-visible noise next to real blocks.
if _is_blank_text_block(b):
return None
out: Dict[str, Any] = _text_block(b["text"])
@@ -339,9 +307,7 @@ def _replay_redacted_thinking(b: Dict[str, Any]) -> Optional[Dict[str, Any]]:
def _replay_tool_use(b: Dict[str, Any]) -> Dict[str, Any]:
out = {
"type": "tool_use",
"id": _sanitize_tool_id(b.get("id", "")),
"name": b.get("name", ""),
"type": "tool_use", "id": _sanitize_tool_id(b.get("id", "")), "name": b.get("name", ""),
"input": b.get("input", {}),
}
if _cache_control_of(b) is not None:
@@ -355,33 +321,24 @@ def _replay_image(b: Dict[str, Any]) -> Optional[Dict[str, Any]]:
_REPLAY_SANITIZERS = {
"text": _replay_text,
"thinking": _replay_thinking,
"redacted_thinking": _replay_redacted_thinking,
"tool_use": _replay_tool_use,
"image": _replay_image,
"text": _replay_text, "thinking": _replay_thinking, "redacted_thinking": _replay_redacted_thinking,
"tool_use": _replay_tool_use, "image": _replay_image,
}
def _sanitize_replay_block(b: Dict[str, Any]) -> Optional[Dict[str, Any]]:
"""Whitelist a stored Anthropic block so it is valid as REQUEST input.
SDK response blocks carry output-only fields the INPUT schema forbids
("Extra inputs are not permitted": ``parsed_output``, ``caller``,
``citations=None``), and ``_to_plain_data`` captured them verbatim. Whitelist
per type (not blacklist) so future SDK fields can't reintroduce the bug;
unknown types are dropped. Returns a clean block or None.
"""
"""Whitelist a stored Anthropic block so it is valid as REQUEST input. SDK response blocks carry
output-only fields the INPUT schema forbids ("Extra inputs are not permitted": ``parsed_output``,
``caller``, ``citations=None``), and ``_to_plain_data`` captured them verbatim. Whitelist per
type (not blacklist) so future SDK fields can't reintroduce the bug; unknown types are dropped.
Returns a clean block or None."""
if not isinstance(b, dict):
return None
sanitizer = _REPLAY_SANITIZERS.get(b.get("type"))
return sanitizer(b) if sanitizer else None
def _apply_assistant_cache_control_to_last_cacheable_block(
blocks: List[Dict[str, Any]],
cache_control: Any,
) -> None:
def _apply_assistant_cache_control_to_last_cacheable_block(blocks: List[Dict[str, Any]], cache_control: Any) -> None:
if not isinstance(cache_control, dict):
return
for block in reversed(blocks):
@@ -390,20 +347,15 @@ def _apply_assistant_cache_control_to_last_cacheable_block(
break
# ---------------------------------------------------------------------------
# Per-message conversion
# ---------------------------------------------------------------------------
# ----- per-message conversion -----
def _replay_ordered_blocks(m: Dict[str, Any], ordered_blocks: List[Any]) -> Optional[List[Dict[str, Any]]]:
"""Interleaved-thinking replay: rebuild the assistant turn from the verbatim
block list normalize_response stored (only for turns interleaving SIGNED
thinking with tool_use). Preserves block ORDER; returns None if nothing survives.
tool_use ``input`` is re-sourced from ``tool_calls`` (redacted at storage time)
rather than the captured block (raw API response, NOT redacted), so a secret
the model inlined into a tool call never goes back on the wire.
"""
"""Interleaved-thinking replay: rebuild the assistant turn from the verbatim block list
normalize_response stored (only for turns interleaving SIGNED thinking with tool_use).
Preserves block ORDER; returns None if nothing survives. tool_use ``input`` is re-sourced from
``tool_calls`` (redacted at storage time) rather than the captured block (raw API response, NOT
redacted), so a secret the model inlined into a tool call never goes back on the wire."""
redacted_input_by_id = {
_sanitize_tool_id(tc.get("id", "")): _parse_tool_args((tc.get("function", {}) or {}).get("arguments", "{}"))
for tc in m.get("tool_calls", []) or []
@@ -425,17 +377,17 @@ def _replay_ordered_blocks(m: Dict[str, Any], ordered_blocks: List[Any]) -> Opti
if redacted is not None:
clean["input"] = redacted
replayed.append(clean)
# Nothing cacheable survived (e.g. signed thinking + blank text): emit the
# placeholder so the turn stays schema-valid and a relocated marker has a carrier.
# Nothing cacheable survived (e.g. signed thinking + blank text): emit the placeholder so the
# turn stays schema-valid and a relocated marker has a carrier.
if not _has_block_type(replayed, _CACHEABLE_TYPES) and (dropped_blank_text or relocated_cc is not None):
replayed.append(_text_block(_EMPTY_TEXT_PLACEHOLDER))
if not replayed:
return None
_apply_assistant_cache_control_to_last_cacheable_block(replayed, relocated_cc)
_apply_assistant_cache_control_to_last_cacheable_block(replayed, m.get("cache_control"))
# prompt_caching marks an assistant turn with text by writing cache_control
# INTO ``content`` (not top-level). This path never reads ``content``, so
# carry that marker over or the breakpoint is burned rather than relocated.
# prompt_caching marks an assistant turn with text by writing cache_control INTO ``content``
# (not top-level). This path never reads ``content``, so carry that marker over or the
# breakpoint is burned rather than relocated.
msg_content = m.get("content")
if isinstance(msg_content, list):
inline_cc = next((cc for cc in map(_cache_control_of, msg_content) if cc is not None), None)
@@ -444,8 +396,8 @@ def _replay_ordered_blocks(m: Dict[str, Any], ordered_blocks: List[Any]) -> Opti
def _convert_assistant_message(m: Dict[str, Any]) -> Dict[str, Any]:
"""Assistant message -> Anthropic content blocks (thinking, text, tool_use,
Kimi/DeepSeek reasoning_content injection)."""
"""Assistant message -> Anthropic content blocks (thinking, text, tool_use, Kimi/DeepSeek
reasoning_content injection)."""
content = m.get("content", "")
ordered_blocks = m.get("anthropic_content_blocks")
if isinstance(ordered_blocks, list) and ordered_blocks:
@@ -454,9 +406,9 @@ def _convert_assistant_message(m: Dict[str, Any]) -> Dict[str, Any]:
return {"role": "assistant", "content": replayed}
blocks = _extract_preserved_thinking_blocks(m)
# Blank text blocks are dropped; a cache marker riding on one is relocated
# onto the last surviving cacheable block (prompt_caching sets cache_control
# on content[-1], which may be exactly the blank block).
# Blank text blocks are dropped; a cache marker riding on one is relocated onto the last
# surviving cacheable block (prompt_caching sets cache_control on content[-1], which may be
# exactly the blank block).
relocated_cc = None
if isinstance(content, list):
for blk in _convert_content_to_anthropic(content):
@@ -472,24 +424,20 @@ def _convert_assistant_message(m: Dict[str, Any]) -> Dict[str, Any]:
continue
fn = tc.get("function", {})
blocks.append({
"type": "tool_use",
"id": _sanitize_tool_id(tc.get("id", "")),
"name": fn.get("name", ""),
"type": "tool_use", "id": _sanitize_tool_id(tc.get("id", "")), "name": fn.get("name", ""),
"input": _parse_tool_args(fn.get("arguments", "{}")),
})
# Kimi's /coding endpoint requires reasoning_content on replayed tool-call
# turns — even "" (injected as a fallback upstream). Prepend, since thinking
# must precede text/tool_use. Skip when reasoning_details already supplied
# (signed) thinking blocks: a duplicate unsigned one would be downgraded to a
# spurious text block on the last assistant message.
# Kimi's /coding endpoint requires reasoning_content on replayed tool-call turns — even ""
# (injected as a fallback upstream). Prepend, since thinking must precede text/tool_use. Skip
# when reasoning_details already supplied (signed) thinking blocks: a duplicate unsigned one
# would be downgraded to a spurious text block on the last assistant message.
reasoning_content = m.get("reasoning_content")
if isinstance(reasoning_content, str) and not _has_block_type(blocks, _THINKING_TYPES):
blocks.insert(0, {"type": "thinking", "thinking": reasoning_content})
# Empty assistant content is rejected. Fall back ONLY to the placeholder,
# never to raw ``content`` — that is the unfiltered blank payload just removed.
# Empty assistant content is rejected. Fall back ONLY to the placeholder, never to raw
# ``content`` — that is the unfiltered blank payload just removed. Markers are applied after
# the fallback so one from a sole dropped blank block lands on the placeholder.
effective = blocks or [_text_block(_EMPTY_TEXT_PLACEHOLDER)]
# Applied after the fallback so a marker from a sole dropped blank block
# lands on the placeholder instead of being lost.
_apply_assistant_cache_control_to_last_cacheable_block(effective, relocated_cc)
_apply_assistant_cache_control_to_last_cacheable_block(effective, m.get("cache_control"))
return {"role": "assistant", "content": effective}
@@ -521,11 +469,10 @@ def _tool_result_content(m: Dict[str, Any]) -> Any:
def _convert_tool_message_to_result(result: List[Dict[str, Any]], m: Dict[str, Any]) -> None:
"""Append a tool_result to ``result``, merging into a trailing tool_result user
message when there is one. Mutates ``result`` in place."""
"""Append a tool_result to ``result``, merging into a trailing tool_result user message when
there is one. Mutates ``result`` in place."""
tool_result = {
"type": "tool_result",
"tool_use_id": _sanitize_tool_id(m.get("tool_call_id", "")),
"type": "tool_result", "tool_use_id": _sanitize_tool_id(m.get("tool_call_id", "")),
"content": _tool_result_content(m),
}
cache_control = _cache_control_of(m)
@@ -552,18 +499,14 @@ def _convert_user_message(content: Any) -> Dict[str, Any]:
return {"role": "user", "content": content}
# ---------------------------------------------------------------------------
# Whole-list passes
# ---------------------------------------------------------------------------
# ----- whole-list passes -----
def _strip_orphaned_tool_blocks(result: List[Dict[str, Any]]) -> None:
"""Strip tool_use blocks with no matching tool_result, and vice versa.
Compression/truncation can remove either side of a pair or insert messages
between them. Anthropic requires the tool_result in the IMMEDIATELY FOLLOWING
user message — a global id match is not enough. Mutates ``result`` in place.
"""
"""Strip tool_use blocks with no matching tool_result, and vice versa. Compression/truncation
can remove either side of a pair or insert messages between them. Anthropic requires the
tool_result in the IMMEDIATELY FOLLOWING user message — a global id match is not enough.
Mutates ``result`` in place."""
# Pass 1: tool_use without an adjacent result.
for i, m in enumerate(result):
if m.get("role") != "assistant" or not isinstance(m.get("content"), list):
@@ -580,9 +523,9 @@ def _strip_orphaned_tool_blocks(result: List[Dict[str, Any]]) -> None:
if not orphaned:
continue
kept = [b for b in m["content"] if not (_block_type(b) == "tool_use" and b.get("id") in orphaned)]
# A signed thinking block on this turn was signed against the ORIGINAL
# content and is now dead (400 "thinking blocks in the latest assistant
# message cannot be modified"). Flag so _manage_thinking_signatures demotes it.
# A signed thinking block on this turn was signed against the ORIGINAL content and is now
# dead (400 "thinking blocks in the latest assistant message cannot be modified"). Flag so
# _manage_thinking_signatures demotes it.
if len(kept) != len(m["content"]) and _has_block_type(m["content"], _THINKING_TYPES):
m["_thinking_signature_invalidated"] = True
m["content"] = kept if kept else [_text_block("(tool call removed)")]
@@ -599,16 +542,15 @@ def _strip_orphaned_tool_blocks(result: List[Dict[str, Any]]) -> None:
if m.get("role") != "user" or not isinstance(m.get("content"), list):
continue
new_content = [
b for b in m["content"]
if _block_type(b) != "tool_result" or b.get("tool_use_id") in surviving_tool_use_ids
b for b in m["content"] if _block_type(b) != "tool_result" or b.get("tool_use_id") in surviving_tool_use_ids
]
if len(new_content) != len(m["content"]):
m["content"] = new_content if new_content else [_text_block("(tool result removed)")]
def _concat_content(prev: Any, curr: Any) -> Any:
"""Merge two message contents: str+str joined by newline, list+list concatenated,
mixed shapes promoted to block lists."""
"""Merge two message contents: str+str joined by newline, list+list concatenated, mixed shapes
promoted to block lists."""
if isinstance(prev, str) and isinstance(curr, str):
return prev + "\n" + curr
if isinstance(prev, str):
@@ -629,24 +571,42 @@ def _merge_consecutive_roles(result: List[Dict[str, Any]]) -> List[Dict[str, Any
# Keep the orphan-strip flag visible to _manage_thinking_signatures.
if m.get("_thinking_signature_invalidated"):
fixed[-1]["_thinking_signature_invalidated"] = True
# The second message's thinking blocks were signed against a
# different turn boundary and become invalid once merged.
# The second message's thinking blocks were signed against a different turn boundary
# and become invalid once merged.
if isinstance(m["content"], list):
m["content"] = [b for b in m["content"] if _block_type(b) not in _THINKING_TYPES]
m["content"] = _strip_thinking(m["content"])
fixed[-1]["content"] = _concat_content(fixed[-1]["content"], m["content"])
return fixed
def _keep_valid_latest_thinking(content: List[Any], signature_dead: bool) -> List[Any]:
"""Latest assistant turn on direct Anthropic: keep signed thinking, demote unsigned to text so
the reasoning isn't lost. If orphan-stripping mutated THIS turn every signature is dead (and a
bare signed block with no tool_use is also invalid), so demote ALL of them."""
new_content = []
for b in content:
if _block_type(b) not in _THINKING_TYPES:
new_content.append(b)
continue
is_redacted = b.get("type") == "redacted_thinking"
signed = b.get("data") if is_redacted else b.get("signature") # redacted 'data' IS the signature
if signed and not signature_dead:
new_content.append(b)
elif (signature_dead or not is_redacted) and b.get("thinking"):
new_content.append(_text_block(b["thinking"])) # demote to plain text
# else: redacted_thinking without data — unverifiable, dropped
return new_content
def _manage_thinking_signatures(result: List[Dict[str, Any]], base_url: str | None, model: str | None) -> None:
"""Strip or preserve thinking blocks per endpoint. Mutates ``result`` in place.
Anthropic signs thinking blocks against the full turn; any upstream mutation
invalidates them (400 "Invalid signature in thinking block"), so on direct
Anthropic only the LATEST assistant turn keeps signed blocks. Signatures are
proprietary: third-party endpoints strip all thinking. Kimi replays as-is;
DeepSeek needs unsigned blocks round-tripped but rejects signed ones. Nous
Portal proxies Claude with sticky sessions and validates the same signatures,
so it takes the native path despite not being anthropic.com.
Anthropic signs thinking blocks against the full turn; any upstream mutation invalidates them
(400 "Invalid signature in thinking block"), so on direct Anthropic only the LATEST assistant
turn keeps signed blocks. Signatures are proprietary: third-party endpoints strip all thinking.
Kimi replays as-is; DeepSeek needs unsigned blocks round-tripped but rejects signed ones. Nous
Portal proxies Claude with sticky sessions and validates the same signatures, so it takes the
native path despite not being anthropic.com.
"""
is_third_party = _is_third_party_anthropic_endpoint(base_url) and not _is_nous_portal_endpoint(base_url)
is_kimi = _is_kimi_family_endpoint(base_url, model)
@@ -666,26 +626,9 @@ def _manage_thinking_signatures(result: List[Dict[str, Any]], base_url: str | No
]
m["content"] = new_content or [_text_block("(empty)")]
elif is_third_party or idx != last_assistant_idx:
stripped = [b for b in m["content"] if _block_type(b) not in _THINKING_TYPES]
m["content"] = stripped or [_text_block("(thinking elided)")]
m["content"] = _strip_thinking(m["content"]) or [_text_block("(thinking elided)")]
else:
# Latest assistant on direct Anthropic: keep signed, demote unsigned to
# text so the reasoning isn't lost. If orphan-stripping mutated THIS
# turn every signature is dead (and a bare signed block with no
# tool_use is also invalid), so demote ALL of them.
signature_dead = bool(m.get("_thinking_signature_invalidated"))
new_content = []
for b in m["content"]:
if _block_type(b) not in _THINKING_TYPES:
new_content.append(b)
continue
is_redacted = b.get("type") == "redacted_thinking"
signed = b.get("data") if is_redacted else b.get("signature") # redacted 'data' IS the signature
if signed and not signature_dead:
new_content.append(b)
elif (signature_dead or not is_redacted) and b.get("thinking"):
new_content.append(_text_block(b["thinking"])) # demote to plain text
# else: redacted_thinking without data — unverifiable, dropped
new_content = _keep_valid_latest_thinking(m["content"], bool(m.get("_thinking_signature_invalidated")))
m["content"] = new_content or [_text_block("(empty)")]
# cache_control on thinking blocks interferes with signature validation.
@@ -696,8 +639,8 @@ def _manage_thinking_signatures(result: List[Dict[str, Any]], base_url: str | No
def _evict_old_screenshots(result: List[Dict[str, Any]]) -> None:
"""Keep only the 3 most recent computer-use screenshots (~1,465 tokens each);
older images become a placeholder text block. Mutates ``result`` in place."""
"""Keep only the 3 most recent computer-use screenshots (~1,465 tokens each); older images
become a placeholder text block. Mutates ``result`` in place."""
image_count = 0
for msg in reversed(result):
content = msg.get("content")
@@ -712,34 +655,25 @@ def _evict_old_screenshots(result: List[Dict[str, Any]]) -> None:
image_count += 1
if image_count > 3:
block["content"] = [
b if b.get("type") != "image" else _text_block("[screenshot removed to save context]")
for b in inner
b if b.get("type") != "image" else _text_block("[screenshot removed to save context]") for b in inner
]
def _ensure_leading_user_turn(result: List[Dict[str, Any]]) -> None:
"""Anthropic requires messages[0].role == user; prepend a placeholder turn otherwise.
A second auto-compaction can leave a role=assistant summary first, which the
API rejects (often masked as a misleading tool_use/tool_result 400). The filler
must be non-whitespace text or it trades that 400 for the blank-block one.
"""
"""Anthropic requires messages[0].role == user; prepend a placeholder turn otherwise. A second
auto-compaction can leave a role=assistant summary first, which the API rejects (often masked
as a misleading tool_use/tool_result 400). The filler must be non-whitespace text or it trades
that 400 for the blank-block one."""
if result and result[0].get("role") != "user":
result.insert(0, {"role": "user", "content": [_text_block(_EMPTY_TEXT_PLACEHOLDER)]})
def _fix_blank_text_blocks_in_list(
blocks: List[Any],
*,
placeholder_text: str,
msg_index: int,
role: Any,
location: str,
blocks: List[Any], *, placeholder_text: str, msg_index: int, role: Any, location: str
) -> List[Any]:
"""Drop blank text blocks; relocate any cache_control they carried onto the last
surviving cacheable block; if nothing survives, substitute one placeholder block
(carrying the relocated marker). Non-text blocks and order are untouched.
Returns a new list; logs structure only (never text)."""
"""Drop blank text blocks; relocate any cache_control they carried onto the last surviving
cacheable block; if nothing survives, substitute one placeholder block (carrying the relocated
marker). Non-text blocks and order are untouched. Returns a new list; logs structure only."""
kept: List[Any] = []
relocated_cache_control = None
for block_index, blk in enumerate(blocks):
@@ -760,12 +694,9 @@ def _fix_blank_text_blocks_in_list(
def _scrub_blank_text_blocks(result: List[Dict[str, Any]]) -> None:
"""Final boundary guard against blank text blocks (HTTP 400 "text content blocks
must contain non-whitespace text"), including inside tool_result content.
Runs LAST so a blank block from any current or future producer never reaches
the wire. Diagnostics are structural only. Mutates ``result`` in place.
"""
"""Final boundary guard against blank text blocks (HTTP 400 "text content blocks must contain
non-whitespace text"), including inside tool_result content. Runs LAST so a blank block from
any current or future producer never reaches the wire. Mutates ``result`` in place."""
for msg_index, msg in enumerate(result):
if not isinstance(msg, dict):
continue
@@ -774,8 +705,7 @@ def _scrub_blank_text_blocks(result: List[Dict[str, Any]]) -> None:
if not isinstance(content, list) or not content:
continue
new_content = _fix_blank_text_blocks_in_list(
content,
placeholder_text=_EMPTY_TEXT_PLACEHOLDER if role == "assistant" else "(empty message)",
content, placeholder_text=_EMPTY_TEXT_PLACEHOLDER if role == "assistant" else "(empty message)",
msg_index=msg_index, role=role, location="content",
)
for blk in new_content:
@@ -790,13 +720,10 @@ def _scrub_blank_text_blocks(result: List[Dict[str, Any]]) -> None:
def _convert_system_content(content: Any) -> Any:
"""System message content -> Anthropic ``system`` param (str, or block list when
cache_control is present).
With cache markers the blocks are copied (never mutating the caller's dicts)
and blank text is replaced by the placeholder: Anthropic rejects blank system
blocks too, and a blank block carrying a breakpoint can't simply be dropped.
"""
"""System message content -> Anthropic ``system`` param (str, or block list when cache_control
is present). With cache markers the blocks are copied (never mutating the caller's dicts) and
blank text is replaced by the placeholder: Anthropic rejects blank system blocks too, and a
blank block carrying a breakpoint can't simply be dropped."""
if not isinstance(content, list):
return content
if not any(p.get("cache_control") for p in content if isinstance(p, dict)):
@@ -812,18 +739,13 @@ def _convert_system_content(content: Any) -> Any:
def convert_messages_to_anthropic(
messages: List[Dict],
base_url: str | None = None,
model: str | None = None,
messages: List[Dict], base_url: str | None = None, model: str | None = None
) -> Tuple[Optional[Any], List[Dict]]:
"""Convert OpenAI-format messages to Anthropic format.
Returns ``(system, messages)``: system is extracted into its own param (a
string, or a block list when cache_control is present). ``base_url``/``model``
drive thinking-signature policy — third-party endpoints strip signatures
"""Convert OpenAI-format messages to Anthropic format -> ``(system, messages)``. System is
extracted into its own param (a string, or a block list when cache_control is present).
``base_url``/``model`` drive thinking-signature policy — third-party endpoints strip signatures
(proprietary, they 400 on them); Kimi-family endpoints/models keep unsigned
reasoning_content-derived blocks, which Kimi requires even when empty.
"""
reasoning_content-derived blocks, which Kimi requires even when empty."""
system = None
result: List[Dict[str, Any]] = []