fix(agent): keep Codex cron calls inline
Cron Codex Responses calls still went through the spawned interrupt worker (the #62151 nested-pool deadlock path). Route them through direct_api_call; the Codex dispatch builds its client via make_client, so the inline stale watchdog aborts a silent Codex stream even though the worker-only TTFB/idle watchdogs are bypassed. Delegated children stay chat_completions-only. Ported onto main's _InlineRequest refactor from PR #70087 (cherry picked from e0028a5d73). Fixes #69734.
This commit is contained in:
@@ -785,13 +785,18 @@ def should_use_direct_api_call(agent) -> bool:
|
||||
thread pools that wedge before the socket opens when the request is pushed onto
|
||||
yet another daemon worker. Running inline drops the deepest layer; interrupts
|
||||
still work because the inline path registers ``agent._active_request_abort``,
|
||||
which ``interrupt()`` invokes cross-thread (#72227). Native/Codex/Bedrock/MoA
|
||||
keep their workers: their cancellation and client ownership differ.
|
||||
which ``interrupt()`` invokes cross-thread (#72227). Cron also inlines Codex
|
||||
Responses (#69734): its dispatch builds the client via ``make_client``, so the
|
||||
inline stale watchdog aborts it like a chat_completions call. Delegated children
|
||||
and Native/Bedrock/MoA keep their workers: cancellation and client ownership differ.
|
||||
"""
|
||||
if getattr(agent, "api_mode", None) != "chat_completions" or getattr(agent, "provider", None) == "moa":
|
||||
api_mode = getattr(agent, "api_mode", None)
|
||||
if getattr(agent, "provider", None) == "moa":
|
||||
return False
|
||||
if getattr(agent, "platform", None) == "cron":
|
||||
return True
|
||||
return api_mode in {"chat_completions", "codex_responses"}
|
||||
if api_mode != "chat_completions":
|
||||
return False
|
||||
# Delegated child — via the execution ContextVar set by _run_single_child,
|
||||
# with the agent's platform stamp as a fallback for callers that bypass it.
|
||||
with contextlib.suppress(Exception):
|
||||
@@ -953,7 +958,7 @@ class _InlineRequest:
|
||||
return newly_stale
|
||||
|
||||
def make_client(self, reason: str, kind: str = "openai"):
|
||||
# Only OpenAI-wire requests reach direct_api_call; ``kind`` exists
|
||||
# Only OpenAI-wire / Codex requests reach direct_api_call; ``kind`` exists
|
||||
# for signature parity with the dispatch helper.
|
||||
client = self.agent._create_request_openai_client(reason=reason, api_kwargs=self.api_kwargs)
|
||||
with self.lock:
|
||||
|
||||
@@ -46,7 +46,11 @@ def test_should_use_direct_api_call_only_for_cron_openai_wire():
|
||||
assert should_use_direct_api_call(_make_agent(platform="telegram")) is False
|
||||
assert should_use_direct_api_call(_make_agent(platform=None)) is False
|
||||
|
||||
for api_mode in ("codex_responses", "anthropic_messages", "bedrock_converse"):
|
||||
codex = _make_agent(platform="cron")
|
||||
codex.api_mode = "codex_responses"
|
||||
assert should_use_direct_api_call(codex) is True # #69734
|
||||
|
||||
for api_mode in ("anthropic_messages", "bedrock_converse"):
|
||||
agent = _make_agent(platform="cron")
|
||||
agent.api_mode = api_mode
|
||||
assert should_use_direct_api_call(agent) is False
|
||||
|
||||
@@ -35,7 +35,11 @@ def test_should_use_direct_api_call_only_for_cron_openai_wire():
|
||||
assert should_use_direct_api_call(_make_agent(platform="telegram")) is False
|
||||
assert should_use_direct_api_call(_make_agent(platform=None)) is False
|
||||
|
||||
for api_mode in ("codex_responses", "anthropic_messages", "bedrock_converse"):
|
||||
codex = _make_agent(platform="cron")
|
||||
codex.api_mode = "codex_responses"
|
||||
assert should_use_direct_api_call(codex) is True # #69734
|
||||
|
||||
for api_mode in ("anthropic_messages", "bedrock_converse"):
|
||||
agent = _make_agent(platform="cron")
|
||||
agent.api_mode = api_mode
|
||||
assert should_use_direct_api_call(agent) is False
|
||||
|
||||
@@ -87,6 +87,29 @@ def test_stalled_inline_call_is_aborted_and_raises_retryable_timeout():
|
||||
|
||||
|
||||
|
||||
def test_stalled_inline_codex_call_is_bounded_by_stale_watchdog():
|
||||
"""#69734: cron Codex runs inline, bypassing the worker-only TTFB/idle
|
||||
watchdogs — the inline stale watchdog must still abort a silent Codex stream."""
|
||||
agent = _make_agent(stale_timeout=0.2)
|
||||
agent.api_mode = "codex_responses"
|
||||
aborted: list[str] = []
|
||||
released = threading.Event()
|
||||
agent._abort_request_openai_client.side_effect = lambda client, reason: (aborted.append(reason), released.set())
|
||||
|
||||
def _silent_codex_stream(api_kwargs, client=None, on_first_delta=None):
|
||||
assert client is agent._create_request_openai_client.return_value
|
||||
if not released.wait(timeout=5.0):
|
||||
raise AssertionError("watchdog never aborted the stalled Codex stream")
|
||||
raise ConnectionError("socket shut down")
|
||||
|
||||
agent._run_codex_stream = _silent_codex_stream
|
||||
started = time.time()
|
||||
with pytest.raises(TimeoutError):
|
||||
direct_api_call(agent, {"model": "gpt-5-codex", "input": []})
|
||||
assert aborted == ["stale_call_kill"]
|
||||
assert time.time() - started < 4.0, "inline watchdog did not bound the Codex call"
|
||||
|
||||
|
||||
def test_watchdog_kill_feeds_the_cross_turn_stale_circuit_breaker():
|
||||
"""Without a bump the #58962 breaker can never trip for cron sessions."""
|
||||
agent = _make_agent(stale_timeout=0.2)
|
||||
|
||||
Reference in New Issue
Block a user