From 10a070bd491087bc724bb907bf5dda5eeb181642 Mon Sep 17 00:00:00 2001 From: kchernev Date: Mon, 24 Aug 2026 10:39:55 +0300 Subject: [PATCH] fix: route codex payloads around the SDK's GIL-holding request transform (#93650) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit responses.create re-walks the entire request body against the ResponseCreateParams union graph client-side while holding the GIL. #93650 documents that walk wedging for 12+ hours on a ~1.4 MB conversation, starving every other thread including the TTFB/stale watchdogs — and no socket kill can unblock a pre-network hang. Hermes payloads are JSON round-trips and already wire format, so the bulk fields (input, tools) are now routed through extra_body, which the SDK merges into the JSON body after the transform. Guarded by a plain-JSON check (anything else keeps the typed path) and a HERMES_CODEX_SDK_TRANSFORM=1 escape hatch. Applied to both the primary stream path and the auxiliary adapter. Co-Authored-By: Claude Fable 5 --- agent/auxiliary_client.py | 8 +- agent/codex_runtime.py | 74 +++++++++ tests/agent/test_auxiliary_client.py | 15 +- .../test_codex_sdk_transform_bypass.py | 153 ++++++++++++++++++ tests/run_agent/test_native_compaction.py | 6 +- .../test_run_agent_codex_responses.py | 7 +- 6 files changed, 259 insertions(+), 4 deletions(-) create mode 100644 tests/run_agent/test_codex_sdk_transform_bypass.py diff --git a/agent/auxiliary_client.py b/agent/auxiliary_client.py index b824a89b3f..251a5e27d2 100644 --- a/agent/auxiliary_client.py +++ b/agent/auxiliary_client.py @@ -1806,10 +1806,16 @@ class _CodexCompletionsAdapter: # Consuming raw events and assembling the final response # ourselves from ``response.output_item.done`` makes us # structurally immune to that drift. - from agent.codex_runtime import _consume_codex_event_stream + from agent.codex_runtime import ( + _bypass_sdk_request_transform, + _consume_codex_event_stream, + ) stream_kwargs = dict(resp_kwargs) stream_kwargs["stream"] = True + # #93650: keep bulk wire-format payload out of the SDK's + # GIL-holding request transform on auxiliary calls too. + stream_kwargs = _bypass_sdk_request_transform(stream_kwargs) def _on_each_event(_event: Any) -> None: # Re-check timeout/cancellation per event, matching the diff --git a/agent/codex_runtime.py b/agent/codex_runtime.py index a09091a635..40e6d05010 100644 --- a/agent/codex_runtime.py +++ b/agent/codex_runtime.py @@ -18,6 +18,7 @@ from __future__ import annotations import json import logging +import os import time from types import SimpleNamespace from typing import Any, Callable, Dict, List @@ -1496,6 +1497,78 @@ def _sanitize_consumer_codex_request( return sanitized +# Bulk request fields that carry the conversation payload. Everything else in +# the request is scalar configuration the SDK transform handles in microseconds. +_SDK_TRANSFORM_BYPASS_FIELDS = ("input", "tools") + + +def _is_plain_json_data(value: Any) -> bool: + """True when ``value`` is composed purely of JSON wire types. + + The SDK's request transform exists to convert typed params (TypedDict + key aliases, pydantic models, ``PropertyInfo`` formats) into wire + format. Hermes assembles Codex payloads from JSON round-trips, so they + are already wire format — but that is only provable when every node is + a plain JSON type. Anything else must keep the typed SDK path. + """ + if value is None or isinstance(value, (str, int, float, bool)): + return True + if isinstance(value, dict): + return all( + isinstance(key, str) and _is_plain_json_data(item) + for key, item in value.items() + ) + if isinstance(value, list): + return all(_is_plain_json_data(item) for item in value) + return False + + +def _bypass_sdk_request_transform(stream_kwargs: dict) -> dict: + """Route bulk payload fields around the SDK's ``maybe_transform`` (#93650). + + ``responses.create`` re-walks the entire request body against the + ``ResponseCreateParams`` union graph before any byte leaves the process. + That walk runs with the GIL held, and #93650 documents it wedging for + 12+ hours on a ~1.4 MB conversation — starving every other thread, + including the TTFB/stale watchdogs whose job is to rescue this exact + call. Because the hang is client-side and pre-network, no socket kill + can unblock it. + + The SDK merges ``extra_body`` into the JSON body *after* the transform + (``_base_client._build_request``), so moving the already-wire-format + bulk fields there skips the walk entirely and produces a byte-identical + request. Fields containing anything that is not plain JSON data (e.g. + pydantic models, generators) stay on the typed path, which still needs + the transform. Set HERMES_CODEX_SDK_TRANSFORM=1 to restore the pre-fix + behavior. + """ + if os.environ.get("HERMES_CODEX_SDK_TRANSFORM", "").strip().lower() in { + "1", "true", "yes", "on" + }: + return stream_kwargs + + moved = { + field: stream_kwargs[field] + for field in _SDK_TRANSFORM_BYPASS_FIELDS + if isinstance(stream_kwargs.get(field), (dict, list)) + and _is_plain_json_data(stream_kwargs[field]) + } + if not moved: + return stream_kwargs + + bypassed = { + key: value for key, value in stream_kwargs.items() if key not in moved + } + extra_body = bypassed.get("extra_body") + merged = dict(extra_body) if isinstance(extra_body, dict) else {} + for field, value in moved.items(): + # An explicit caller-provided extra_body entry keeps precedence, + # matching what the SDK's post-transform merge would have done. + merged.setdefault(field, value) + bypassed["extra_body"] = merged + return bypassed + + def run_codex_stream(agent, api_kwargs: dict, client: Any = None, on_first_delta=None): """Execute one streaming Responses API request and return the final response. @@ -1543,6 +1616,7 @@ def run_codex_stream(agent, api_kwargs: dict, client: Any = None, on_first_delta next_api_kwargs, ) stream_kwargs["stream"] = True + stream_kwargs = _bypass_sdk_request_transform(stream_kwargs) return active_client.responses.create(**stream_kwargs) def _codex_stream_created(_raw_stream: Any) -> None: diff --git a/tests/agent/test_auxiliary_client.py b/tests/agent/test_auxiliary_client.py index 238c9e5038..1d784efa1e 100644 --- a/tests/agent/test_auxiliary_client.py +++ b/tests/agent/test_auxiliary_client.py @@ -2830,6 +2830,9 @@ class TestCodexAdapterReasoningTranslation: def _create(**kwargs): captured_kwargs.update(kwargs) + # #93650 routes bulk fields through extra_body; fold them back in + # so assertions read the effective wire body the SDK would send. + captured_kwargs.update(kwargs.get("extra_body") or {}) return _FakeCreateStream() real_client = MagicMock() @@ -2913,6 +2916,9 @@ class TestCodexAdapterPromptCacheKey: def _create(**kwargs): captured_kwargs.update(kwargs) + # #93650 routes bulk fields through extra_body; fold them back in + # so assertions read the effective wire body the SDK would send. + captured_kwargs.update(kwargs.get("extra_body") or {}) return _FakeCreateStream() real_client = MagicMock() @@ -3028,6 +3034,9 @@ class TestCodexAdapterGithubResponsesMessageIdDrop: def _create(**kwargs): captured_kwargs.update(kwargs) + # #93650 routes bulk fields through extra_body; fold them back in + # so assertions read the effective wire body the SDK would send. + captured_kwargs.update(kwargs.get("extra_body") or {}) return _FakeCreateStream() real_client = MagicMock() @@ -3312,7 +3321,11 @@ class TestCodexAuxiliaryToolMessageConversion: fake_client = SimpleNamespace(responses=FakeResponses()) adapter = _CodexCompletionsAdapter(fake_client, "gpt-5.5") adapter.create(messages=messages, model="gpt-5.5") - return fake_client.responses.kwargs + # #93650 routes bulk fields through extra_body; fold them back in so + # assertions read the effective wire body the SDK would send. + kwargs = dict(fake_client.responses.kwargs) + kwargs.update(kwargs.pop("extra_body", None) or {}) + return kwargs def test_tool_history_never_leaks_role_tool(self): messages = [ diff --git a/tests/run_agent/test_codex_sdk_transform_bypass.py b/tests/run_agent/test_codex_sdk_transform_bypass.py new file mode 100644 index 0000000000..ce9187779c --- /dev/null +++ b/tests/run_agent/test_codex_sdk_transform_bypass.py @@ -0,0 +1,153 @@ +"""Regression tests for the SDK request-transform bypass (#93650). + +``responses.create`` re-walks the whole request body against the +``ResponseCreateParams`` union graph client-side, holding the GIL. #93650 +documents that walk wedging for 12+ hours on a ~1.4 MB conversation and +freezing the entire agent — no in-process watchdog can fire while the GIL +is held, and no socket kill helps a pre-network hang. Bulk wire-format +fields are therefore routed through ``extra_body``, which the SDK merges +into the JSON body *after* the transform. +""" +import sys +import types +from types import SimpleNamespace +from unittest.mock import MagicMock + +sys.modules.setdefault("fire", types.SimpleNamespace(Fire=lambda *a, **k: None)) +sys.modules.setdefault("firecrawl", types.SimpleNamespace(Firecrawl=object)) +sys.modules.setdefault("fal_client", types.SimpleNamespace()) + +from agent.codex_runtime import ( + _bypass_sdk_request_transform, + _is_plain_json_data, +) + + +def _wire_kwargs(): + return { + "model": "gpt-5.6-sol", + "instructions": "You are Hermes.", + "input": [ + {"role": "user", "content": [{"type": "input_text", "text": "Ping"}]}, + {"type": "function_call_output", "call_id": "c1", "output": "ok"}, + ], + "tools": [{"type": "function", "name": "terminal", "parameters": {}}], + "store": False, + "stream": True, + "timeout": 1800.0, + } + + +class TestIsPlainJsonData: + def test_accepts_nested_wire_payloads(self): + assert _is_plain_json_data(_wire_kwargs()["input"]) + + def test_rejects_non_json_leaves(self): + assert not _is_plain_json_data([{"role": "user", "content": object()}]) + + def test_rejects_non_string_dict_keys(self): + assert not _is_plain_json_data({1: "a"}) + + def test_rejects_generators(self): + assert not _is_plain_json_data((item for item in ())) + + +class TestBypassSdkRequestTransform: + def test_moves_bulk_fields_to_extra_body(self): + kwargs = _wire_kwargs() + original_input = kwargs["input"] + + bypassed = _bypass_sdk_request_transform(kwargs) + + assert "input" not in bypassed + assert "tools" not in bypassed + assert bypassed["extra_body"]["input"] is original_input + assert bypassed["extra_body"]["tools"] == kwargs["tools"] + # Scalar configuration stays on the typed path. + assert bypassed["model"] == "gpt-5.6-sol" + assert bypassed["stream"] is True + assert bypassed["timeout"] == 1800.0 + # The caller's mapping is untouched. + assert kwargs["input"] is original_input + assert "extra_body" not in kwargs + + def test_merges_with_existing_extra_body_and_keeps_caller_precedence(self): + kwargs = _wire_kwargs() + caller_extra = {"prompt_cache_retention": "24h", "input": "explicit-wins"} + kwargs["extra_body"] = caller_extra + + bypassed = _bypass_sdk_request_transform(kwargs) + + # An explicit extra_body entry wins, exactly as the SDK's + # post-transform merge would have resolved the collision. + assert bypassed["extra_body"]["input"] == "explicit-wins" + assert bypassed["extra_body"]["prompt_cache_retention"] == "24h" + assert bypassed["extra_body"]["tools"] == kwargs["tools"] + assert caller_extra == { + "prompt_cache_retention": "24h", + "input": "explicit-wins", + } + + def test_non_json_field_stays_on_typed_sdk_path(self): + kwargs = _wire_kwargs() + kwargs["input"] = [{"role": "user", "content": object()}] + + bypassed = _bypass_sdk_request_transform(kwargs) + + assert bypassed["input"] == kwargs["input"] + assert bypassed["extra_body"] == {"tools": kwargs["tools"]} + + def test_string_input_stays_in_place(self): + kwargs = _wire_kwargs() + kwargs["input"] = "plain prompt" + kwargs.pop("tools") + + bypassed = _bypass_sdk_request_transform(kwargs) + + assert bypassed is kwargs + + def test_env_escape_hatch_restores_passthrough(self, monkeypatch): + monkeypatch.setenv("HERMES_CODEX_SDK_TRANSFORM", "1") + kwargs = _wire_kwargs() + + assert _bypass_sdk_request_transform(kwargs) is kwargs + + +class TestRunCodexStreamRoutesPayloadViaExtraBody: + def _make_agent(self): + from run_agent import AIAgent + + agent = AIAgent( + api_key="test-key", + base_url="https://chatgpt.com/backend-api/codex", + model="gpt-5.6-sol", + quiet_mode=True, + skip_context_files=True, + skip_memory=True, + ) + agent._interrupt_requested = False + return agent + + def test_create_receives_input_via_extra_body(self): + from agent.codex_runtime import run_codex_stream + + agent = self._make_agent() + events = [ + SimpleNamespace( + type="response.completed", + response=SimpleNamespace( + id="r1", status="completed", output=[], usage=None, + ), + ) + ] + mock_client = MagicMock() + mock_client.responses.create.return_value = iter(events) + + run_codex_stream(agent, _wire_kwargs(), client=mock_client) + + create_kwargs = mock_client.responses.create.call_args.kwargs + assert "input" not in create_kwargs + assert "tools" not in create_kwargs + assert create_kwargs["stream"] is True + assert create_kwargs["extra_body"]["input"][0]["role"] == "user" + assert create_kwargs["extra_body"]["tools"][0]["name"] == "terminal" diff --git a/tests/run_agent/test_native_compaction.py b/tests/run_agent/test_native_compaction.py index 6518b44fbf..bf7150c32e 100644 --- a/tests/run_agent/test_native_compaction.py +++ b/tests/run_agent/test_native_compaction.py @@ -694,7 +694,11 @@ class TestCheckpointGatedOnCurrentEligibility: class _Responses: def create(self, **kwargs): - seen["input"] = kwargs.get("input") + # #93650 routes the bulk input around the SDK transform via + # extra_body; accept the payload in either wire shape. + seen["input"] = kwargs.get("input") or ( + kwargs.get("extra_body") or {} + ).get("input") raise RuntimeError("stop before network") class _Client: diff --git a/tests/run_agent/test_run_agent_codex_responses.py b/tests/run_agent/test_run_agent_codex_responses.py index 97e1ed1679..3e377c86f5 100644 --- a/tests/run_agent/test_run_agent_codex_responses.py +++ b/tests/run_agent/test_run_agent_codex_responses.py @@ -630,7 +630,12 @@ def test_run_codex_stream_strips_nested_request_override_retention( with caplog.at_level("WARNING", logger="agent.codex_runtime"): agent._run_codex_stream(request) - assert "extra_body" not in captured + # The transform bypass (#93650) re-introduces extra_body to carry the bulk + # payload fields; the guard's contract is that retention itself never + # crosses the wire boundary in either shape. + assert "prompt_cache_retention" not in captured + assert "prompt_cache_retention" not in captured.get("extra_body", {}) + assert "input" in captured.get("extra_body", {}) assert request["extra_body"] == {"prompt_cache_retention": "24h"} assert any( "Dropped unsupported prompt_cache_retention at consumer Codex wire boundary"