From 1cd2ed6c20d68df9dde324640fa3470a77b71c79 Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Wed, 2 Sep 2026 23:18:46 -0700 Subject: [PATCH] refactor(agent/codex_responses_adapter): merge image-part helpers, alias _deterministic_call_id, compact docstrings/comments --- agent/codex_responses_adapter.py | 104 +++++++++++-------------------- 1 file changed, 38 insertions(+), 66 deletions(-) diff --git a/agent/codex_responses_adapter.py b/agent/codex_responses_adapter.py index f3048158ac..7559841d3a 100644 --- a/agent/codex_responses_adapter.py +++ b/agent/codex_responses_adapter.py @@ -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.