refactor(agent/codex_responses_adapter): merge image-part helpers, alias _deterministic_call_id, compact docstrings/comments

This commit is contained in:
Teknium
2026-09-02 23:18:46 -07:00
parent bb1b6d4a36
commit 1cd2ed6c20
+38 -66
View File
@@ -1,8 +1,5 @@
"""Codex Responses API adapter.
Stateless format-conversion and normalization for the OpenAI Responses API
(OpenAI Codex, xAI, GitHub Models and other Responses-compatible endpoints).
"""
"""Codex Responses API adapter: stateless format conversion and normalization for the
OpenAI Responses API (OpenAI Codex, xAI, GitHub Models and other compatible endpoints)."""
from __future__ import annotations
@@ -60,8 +57,8 @@ _VALID_RESPONSES_FN_NAME_RE = re.compile(r"[a-zA-Z0-9_-]{1,64}")
# Provider-executed built-in tools: declared by ``type`` alone, run server-side,
# reported via the ``*_call`` output items below; preflight passes them through.
_RESPONSES_BUILTIN_TOOL_TYPES = {
"web_search", "web_search_preview", "file_search", "code_interpreter",
"image_generation", "computer_use_preview", "local_shell",
"web_search", "web_search_preview", "file_search", "code_interpreter", "image_generation", "computer_use_preview",
"local_shell",
}
# Server-side ``*_call`` output items. xAI leaves these ``in_progress`` even when the
@@ -119,9 +116,8 @@ def _neutralize_harmony_tokens(text: str) -> str:
"""Keep Harmony source readable without emitting reserved wire tokens."""
if not text or "<" not in text or "|" not in text:
return text
replacement = rf"<{_FULLWIDTH_PIPE}\1{_FULLWIDTH_PIPE}>"
if not any(unicodedata.category(char) == "Cf" for char in text):
return _HARMONY_CONTROL_TOKEN_RE.sub(replacement, text)
return _HARMONY_CONTROL_TOKEN_RE.sub(rf"<{_FULLWIDTH_PIPE}\1{_FULLWIDTH_PIPE}>", text)
# The backend strips Unicode format controls (e.g. U+200B) before its reserved-token
# check, so match on the visible text and rewrite the original spans.
original_positions = [i for i, char in enumerate(text) if unicodedata.category(char) != "Cf"]
@@ -168,9 +164,12 @@ def _iter_content_parts(content: list) -> Iterator[tuple[str, Any]]:
yield "image", part
def _input_image_part(part: Dict[str, Any], *, keep_empty_url: bool) -> Optional[Dict[str, Any]]:
"""``input_image`` from a chat/Responses image part (``image_url`` may be a str or ``{url, detail}``);
None for an empty url unless ``keep_empty_url``."""
def _input_image_part(part: Dict[str, Any], role: str = "user", *, keep_empty_url: bool) -> Optional[Dict[str, Any]]:
"""Responses image part from a chat/Responses image part (``image_url`` may be a str or
``{url, detail}``). Assistant → text placeholder (an assistant ``input_image`` 400s every
replay); user → ``input_image``, None for an empty url unless ``keep_empty_url``."""
if role == "assistant":
return {"type": "output_text", "text": _ASSISTANT_IMAGE_PLACEHOLDER}
url, detail = part.get("image_url"), part.get("detail")
if isinstance(url, dict):
url, detail = url.get("url"), url.get("detail", detail)
@@ -182,26 +181,16 @@ def _input_image_part(part: Dict[str, Any], *, keep_empty_url: bool) -> Optional
return image_part
def _image_part_for_role(part: Dict[str, Any], role: str, *, keep_empty_url: bool) -> Optional[Dict[str, Any]]:
"""Responses image part for ``role``: assistant → text placeholder (an assistant
``input_image`` 400s every replay); user → ``input_image``."""
if role == "assistant":
return {"type": "output_text", "text": _ASSISTANT_IMAGE_PLACEHOLDER}
return _input_image_part(part, keep_empty_url=keep_empty_url)
def _chat_content_to_responses_parts(content: Any, *, role: str = "user") -> List[Dict[str, Any]]:
"""Chat-style multimodal content → Responses API input parts ([] if not a list).
Text becomes ``input_text`` (user) or ``output_text`` (assistant) — the API rejects
the wrong type per role. ``input_image`` is only legal on user messages; on assistant
messages it becomes a text marker (an assistant ``input_image`` 400s every replay)."""
"""Chat-style multimodal content → Responses API input parts ([] if not a list). Text is
``input_text`` (user) / ``output_text`` (assistant) — the API rejects the wrong type per role;
``input_image`` is only legal on user messages (see :func:`_input_image_part`)."""
text_type = _text_type_for(role)
converted: List[Dict[str, Any]] = []
for kind, payload in _iter_content_parts(_as_list(content)):
if kind == "text":
converted.append({"type": text_type, "text": payload})
elif (part := _image_part_for_role(payload, role, keep_empty_url=False)) is not None:
elif (part := _input_image_part(payload, role, keep_empty_url=False)) is not None:
converted.append(part)
return converted
@@ -210,8 +199,6 @@ def _summarize_user_message_for_log(content: Any, *, sep: str = " ") -> str:
"""Flatten message content to plain text: text parts joined with ``sep`` (``" "`` for
logs; ``"\\n"`` for memory providers feeding regexes), images → ``[N image(s)]``
marker, ``""`` for None/empty, ``str(content)`` for other scalars."""
if isinstance(content, str):
return content
if not isinstance(content, list):
try:
return _str_or_empty(content)
@@ -226,9 +213,8 @@ def _summarize_user_message_for_log(content: Any, *, sep: str = " ") -> str:
# --- ID helpers ---------------------------------------------------------------
def _deterministic_call_id(fn_name: str, arguments: str, index: int = 0) -> str:
"""Deterministic call_id (random ids would break the prompt-cache prefix). Re-exported for run_agent/tests."""
return deterministic_call_id(fn_name, arguments, index)
# Deterministic call_id fallback (random ids would break the prompt-cache prefix); re-exported for run_agent/tests.
_deterministic_call_id = deterministic_call_id
def _clamp_responses_call_id(call_id: str) -> str:
@@ -241,10 +227,9 @@ def _clamp_responses_call_id(call_id: str) -> str:
def _sanitize_replayed_fn_name(name: str) -> str:
"""Coerce a *replayed* ``function_call.name`` to ``^[a-zA-Z0-9_-]{1,64}$`` (an invalid
stored name 400s every later turn). Invalid runs collapse to ``_``; all-invalid → "fn".
Apply ONLY to replayed items, never live tool definitions (schema names must match the
dispatch registry); pairing with the output is by call_id."""
"""Coerce a *replayed* ``function_call.name`` to ``^[a-zA-Z0-9_-]{1,64}$`` (an invalid stored
name 400s every later turn). Invalid runs collapse to ``_``; all-invalid → "fn". Apply ONLY to
replayed items, never live tool definitions (schema names must match the dispatch registry)."""
if not isinstance(name, str):
return "fn"
if _VALID_RESPONSES_FN_NAME_RE.fullmatch(name):
@@ -340,9 +325,7 @@ def _message_item(
return item
def _assistant_message_item(
raw: Dict[str, Any], content: List[Dict[str, Any]], *, is_github_responses: bool,
) -> Dict[str, Any]:
def _assistant_message_item(raw: Dict[str, Any], content: List[Dict[str, Any]], *, is_github_responses: bool) -> Dict[str, Any]:
"""Replayable assistant ``message`` item from a stored one. ``id`` is kept only when
short enough and never for GitHub Copilot (ids bind to a backend connection; stale →
401); ``phase`` is preserved per OpenAI's cache guidance."""
@@ -357,12 +340,11 @@ def _assistant_message_item(
def _replay_reasoning_items(
msg: Dict[str, Any], *, seen_item_ids: set, current_issuer_kind: Optional[str], native_compaction_eligible: bool,
) -> List[Dict[str, Any]]:
"""Replay persisted encrypted reasoning/compaction items for one assistant turn.
Skips duplicate ids, ``compaction`` checkpoints unless THIS request carries
``context_management`` (else a persisted checkpoint erases pre-checkpoint history on a
model that cannot decrypt it), and items stamped by another issuer (HTTP 400).
Unstamped legacy items pass. ``id`` (store=False lookups 404) and ``_issuer_kind`` are stripped."""
"""Replay persisted encrypted reasoning/compaction items for one assistant turn. Skips duplicate
ids, ``compaction`` checkpoints unless THIS request carries ``context_management`` (else a persisted
checkpoint erases pre-checkpoint history on a model that cannot decrypt it), and items stamped by
another issuer (HTTP 400); unstamped legacy items pass. ``id`` (store=False lookups 404) and
``_issuer_kind`` are stripped."""
global _CROSS_ISSUER_WARN_EMITTED
replayed: List[Dict[str, Any]] = []
for ri in _as_list(msg.get("codex_reasoning_items")):
@@ -456,17 +438,15 @@ def _chat_messages_to_responses_input(
``AIAgent._disable_codex_reasoning_replay`` after an ``invalid_encrypted_content`` 400.
``is_github_responses``: drops ``id`` from replayed message items (Copilot 401s on stale ids).
``current_issuer_kind``: cross-issuer guard; foreign-stamped items drop, legacy items replay.
``native_compaction_eligible``: THIS request carries ``context_management``; gates both
replaying ``compaction`` checkpoints and ``prune_pre_checkpoint_items``. Checkpoints
persist across model swaps / compression flips / resume, so without the gate one
checkpoint would delete pre-checkpoint history on a model that cannot decrypt it.
Dropping it is lossless: local history is never truncated by native compaction."""
``native_compaction_eligible``: THIS request carries ``context_management``; gates both replaying
``compaction`` checkpoints and ``prune_pre_checkpoint_items``. Checkpoints persist across model
swaps / compression flips / resume, so without the gate one checkpoint would delete pre-checkpoint
history on a model that cannot decrypt it. Dropping it is lossless: local history is never truncated."""
items: List[Dict[str, Any]] = []
# Parallel to ``items``: source chat message per item. Pruning reads a summary
# carrier's provenance from the source; the converted item may be a lossy shape.
item_sources: List[Optional[Dict[str, Any]]] = []
seen_item_ids: set = set()
def emit(new_items: List[Dict[str, Any]], msg: Dict[str, Any]) -> None:
items.extend(new_items)
item_sources.extend([msg] * len(new_items))
@@ -502,10 +482,9 @@ def _chat_messages_to_responses_input(
if fallback is not None:
emit([{"role": "assistant", "content": fallback}], msg)
emit(_replay_tool_call_items(msg, start_index=len(items)), msg)
# The server renders nothing placed before a compaction item, so pre-checkpoint
# history is dead weight and plaintext asks / merged summaries silently vanish. Keep
# the newest checkpoint first, retain pre-checkpoint USER and SUMMARY messages within
# a token budget, leave the tail untouched.
# The server renders nothing placed before a compaction item, so pre-checkpoint history is
# dead weight and plaintext asks / merged summaries silently vanish. Keep the newest checkpoint
# first, retain pre-checkpoint USER and SUMMARY messages within a token budget, leave the tail.
if not native_compaction_eligible:
return items
from agent.native_compaction import prune_pre_checkpoint_items
@@ -515,7 +494,6 @@ def _chat_messages_to_responses_input(
class ResponsesRouteFlags(NamedTuple):
"""Which special Responses-API route an agent is talking to. Single owner of the
codex/xai/github predicates — every site must call :func:`classify_responses_route`."""
is_codex_backend: bool
is_xai_responses: bool
is_github_responses: bool
@@ -529,7 +507,6 @@ def classify_responses_route(agent: Any) -> ResponsesRouteFlags:
base_url = str(getattr(agent, "base_url", "") or "")
hostname = str(getattr(agent, "_base_url_hostname", "") or "").lower() or base_url_hostname(base_url)
lower = str(getattr(agent, "_base_url_lower", "") or base_url).lower()
def _host_is(domain: str) -> bool:
return hostname == domain or hostname.endswith("." + domain)
return ResponsesRouteFlags(
@@ -672,7 +649,7 @@ def _preflight_role_message(item: Dict[str, Any], idx: int, ctx: _PreflightCtx)
text = text if isinstance(text, str) else str(text or "")
validated.append({"type": text_type, "text": ctx.sanitize_text(text)})
elif ptype in _IMAGE_PART_TYPES:
validated.append(_image_part_for_role(part, role, keep_empty_url=True))
validated.append(_input_image_part(part, role, keep_empty_url=True))
else:
raise ValueError(
f"Codex Responses input[{idx}].content[{part_idx}] has unsupported type {part.get('type')!r}."
@@ -936,9 +913,7 @@ class _OutputScan:
logger.info(
"Native Responses compaction item captured (%d chars encrypted).", len(raw_item["encrypted_content"]),
)
elif item_type in {"function_call", "custom_tool_call"}:
if item_type == "function_call" and item_status in _INCOMPLETE_STATUSES:
continue
elif item_type == "custom_tool_call" or (item_type == "function_call" and item_status not in _INCOMPLETE_STATUSES):
self.tool_calls.append(_response_tool_call(item, item_type, len(self.tool_calls)))
def _message(self, item: Any, item_status: Optional[str]) -> None:
@@ -954,8 +929,7 @@ class _OutputScan:
(self.reasoning_parts if is_commentary_phase else self.content_parts).append(message_text)
item_id = getattr(item, "id", None)
self.message_items_raw.append(_message_item(
[{"type": "output_text", "text": message_text}],
status=_normalize_responses_message_status(item_status),
[{"type": "output_text", "text": message_text}], status=_normalize_responses_message_status(item_status),
item_id=item_id if isinstance(item_id, str) else None, phase=normalized_phase,
))
@@ -964,10 +938,8 @@ def _normalize_codex_response(response: Any, *, issuer_kind: Optional[str] = Non
"""Normalize a Responses API object to ``(assistant_message, finish_reason)``.
``issuer_kind`` is stamped onto captured reasoning items for cross-issuer replay drops."""
response_status = _lower_or_none(getattr(response, "status", None))
incomplete_reason = _field(getattr(response, "incomplete_details", None), "reason", "")
response_incomplete_content_filter = (
response_status == "incomplete" and str(incomplete_reason or "").strip().lower() == "content_filter"
)
incomplete_reason = str(_field(getattr(response, "incomplete_details", None), "reason", "") or "").strip().lower()
response_incomplete_content_filter = response_status == "incomplete" and incomplete_reason == "content_filter"
output = getattr(response, "output", None)
if not isinstance(output, list) or not output:
# Codex can deliver the whole answer via stream events with an empty output.