From 5eb0f4f7d125020400c76dabaaf52b2033d4677c Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Wed, 2 Sep 2026 19:03:50 -0700 Subject: [PATCH] refactor(agent/anthropic_message_convert): share _strip_thinking, extract latest-turn thinking filter, compact docstrings; endpoints docstrings --- agent/anthropic_endpoints.py | 110 +++------ agent/anthropic_message_convert.py | 376 ++++++++++++----------------- 2 files changed, 188 insertions(+), 298 deletions(-) diff --git a/agent/anthropic_endpoints.py b/agent/anthropic_endpoints.py index a9052567d2..fbd871dde0 100644 --- a/agent/anthropic_endpoints.py +++ b/agent/anthropic_endpoints.py @@ -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 diff --git a/agent/anthropic_message_convert.py b/agent/anthropic_message_convert.py index f3bca75ff6..c053be4311 100644 --- a/agent/anthropic_message_convert.py +++ b/agent/anthropic_message_convert.py @@ -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]] = []