From 9387bf929c39bf31ffbcbdb9f7f97c370ff9d8b8 Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Tue, 1 Sep 2026 11:30:27 -0700 Subject: [PATCH] fix(delegate): drain abandoned-worker transports FD-safely on child timeout MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The #94248 native half. A delegation deadline abandons the child's daemon worker while it is typically parked inside an in-flight OpenSSL read (Codex Responses stream / httpx). PR #90889's deferred close (cherry-picked here, authorship preserved) stops the timeout thread from closing the child under the running future — but the deferred close only fires once the worker unwinds, and a worker blocked in ssl.read never unwinds on its own: the cooperative interrupt cannot reach a thread inside OpenSSL, so the child's SessionDB, httpx pools, and subprocesses stayed pinned until process exit, and any path that still hard-closed the transport released FDs under a live SSL BIO (the #29507/#67142/#70773 native-corruption family; SIGSEGV 17-72ms after "Subagent N timed out" on macOS arm64). Fix — bounded drain after deferral: - AIAgent._drain_transports_after_abandonment(): shutdown()-only sweep of the shared client's pooled sockets (force_close_tcp_sockets — FD release stays with the owning worker), abort+poison of the cached per-request openai/anthropic wire clients, Codex app-server request_interrupt(), and the inline _active_request_abort hook. Never client.close(), never socket.close(). - delegate timeout path: after registering the deferred-close callback, run one immediate drain plus one 5s re-sweep (covers a connection opened between the interrupt and the first sweep). The settled read (EOF/EPIPE) lets the worker unwind, which triggers the deferred close on the worker's own thread — the only safe FD-release boundary. A worker that still never settles retains its resources rather than risking a cross-thread close. Live repro (Linux, real TLS server subprocess + real httpx client blocked in OpenSSL read at the deadline + real SessionDB): before — child.close() ran on the timeout thread with in_flight_ssl_read=True (client FDs released under the live read; #94736 self-heal WARNING fired on the worker's unwind flush); after — drain settles the read in ~1ms, worker unwinds, close runs on the worker thread with in_flight_ssl_read=False. Not live-tested on macOS arm64 (no macOS runner); the fix is platform-neutral teardown ordering proven on Linux. Closes #94248 --- run_agent.py | 69 ++++++ .../test_94248_timeout_transport_drain.py | 197 ++++++++++++++++++ tools/delegate_tool.py | 35 ++++ 3 files changed, 301 insertions(+) create mode 100644 tests/tools/test_94248_timeout_transport_drain.py diff --git a/run_agent.py b/run_agent.py index 261945efd2..2748dcfa61 100644 --- a/run_agent.py +++ b/run_agent.py @@ -5599,6 +5599,75 @@ class AIAgent: exc, ) + def _drain_transports_after_abandonment(self, *, reason: str) -> int: + """FD-safe transport drain for an abandoned (timed-out) worker (#94248). + + A delegation deadline abandons this agent's daemon worker while it may + still be blocked inside an in-flight OpenSSL ``read`` (Codex Responses + stream, httpx request). The timeout thread must never hard-close those + transports — ``client.close()`` releases raw FDs under a live SSL BIO, + the #29507 / #67142 / #70773 native-corruption family and the SIGSEGV + shape reported in #94248. This helper only ``shutdown()``s pooled + sockets (safe from any thread), settling blocked reads with EOF/EPIPE + so the worker can unwind and run the real close from its own thread. + + Returns the number of sockets shut down across all transports. + """ + drained = 0 + # Shared primary client (codex-direct / MoA stream on it directly). + try: + client = getattr(self, "client", None) + if client is not None: + drained += self._force_close_tcp_sockets(client) + except Exception: + logger.debug("Abandoned-worker drain: shared client sweep failed", + exc_info=True) + # Cached per-request wire clients: abort (shutdown + poison the reuse + # slot) so the unwinding worker discards them instead of re-caching. + try: + with self._openai_client_lock(): + cache = getattr(self, "_request_client_cache", None) + cached = cache["client"] if cache else None + if cached is not None: + self._abort_request_openai_client(cached, reason=reason) + except Exception: + logger.debug("Abandoned-worker drain: request client abort failed", + exc_info=True) + try: + with self._openai_client_lock(): + cache = getattr(self, "_request_anthropic_client_cache", None) + cached = cache["client"] if cache else None + if cached is not None: + self._abort_request_anthropic_client(cached, reason=reason) + except Exception: + logger.debug("Abandoned-worker drain: anthropic client abort failed", + exc_info=True) + # Codex app-server session watches a private interrupt event. + try: + codex_session = getattr(self, "_codex_session", None) + request_interrupt = getattr(codex_session, "request_interrupt", None) + if callable(request_interrupt): + request_interrupt() + except Exception: + logger.debug("Abandoned-worker drain: codex interrupt failed", + exc_info=True) + # Inline (cron-style) request abort hook, when registered. + try: + abort_active = getattr(self, "_active_request_abort", None) + if callable(abort_active): + abort_active(reason) + except Exception: + logger.debug("Abandoned-worker drain: active request abort failed", + exc_info=True) + logger.info( + "Abandoned-worker transports drained (%s, tcp_shutdown=%d, " + "fd_release=deferred_to_worker) %s", + reason, + drained, + self._client_log_context(), + ) + return drained + def _build_primary_client_for_active_provider(self, *, reason: str) -> Any: """Build the shared client shape required by the active provider. diff --git a/tests/tools/test_94248_timeout_transport_drain.py b/tests/tools/test_94248_timeout_transport_drain.py new file mode 100644 index 0000000000..2347ad2915 --- /dev/null +++ b/tests/tools/test_94248_timeout_transport_drain.py @@ -0,0 +1,197 @@ +"""#94248 (native half): delegation timeout must drain transports FD-safely. + +A timed-out child's daemon worker is typically parked inside an in-flight +OpenSSL read. The timeout thread must (1) never hard-close the child while the +worker future is running (deferred close, #90889), and (2) drain the child's +transports with socket ``shutdown()`` only — never ``client.close()`` — so the +blocked read settles with EOF/EPIPE and the worker can unwind (bounded drain). +Cross-thread FD release under a live SSL BIO is the #29507/#67142/#70773 +native-corruption family. +""" +from __future__ import annotations + +import threading +import time +from types import SimpleNamespace + +from tools import delegate_tool + + +class _SslBlockedChild: + """Worker blocks (modelling an in-flight SSL read) until drained.""" + + def __init__(self) -> None: + self.tool_progress_callback = None + self._credential_pool = None + self._delegate_saved_tool_names = [] + self._delegate_role = "leaf" + self._delegate_depth = 1 + self._subagent_id = None + self.model = "test-model" + self.session_prompt_tokens = 0 + self.session_completion_tokens = 0 + self.session_estimated_cost_usd = 0.0 + self.session_cost_status = "unknown" + self.read_settled = threading.Event() # drain "EOF" signal + self.unwound = threading.Event() + self.closed = threading.Event() + self.close_while_blocked = False + self.drain_calls: list[str] = [] + self.drain_threads: list[str] = [] + + def run_conversation(self, **_kwargs): + # Models the worker blocked in ssl.read: only the FD-safe drain + # (socket shutdown -> EOF) settles it; interrupts alone do not. + assert self.read_settled.wait(timeout=10), "drain never settled the read" + time.sleep(0.05) # post-read unwind work (turn-finally flush) + self.unwound.set() + return { + "final_response": "", + "completed": False, + "interrupted": True, + "api_calls": 1, + "messages": [], + } + + def hard_interrupt(self, *_a, **_k): + # Cooperative interrupt cannot unblock a thread inside OpenSSL read. + pass + + def get_activity_summary(self): + return {"api_call_count": 1} + + def _drain_transports_after_abandonment(self, *, reason: str) -> int: + self.drain_calls.append(reason) + self.drain_threads.append(threading.current_thread().name) + self.read_settled.set() + return 1 + + def close(self): + if not self.unwound.is_set(): + self.close_while_blocked = True + self.closed.set() + + +def _run(child, monkeypatch, timeout=0.4): + parent = SimpleNamespace( + session_id="parent-94248-drain", + _current_task_id=None, + _active_children=[child], + _active_children_lock=threading.Lock(), + ) + monkeypatch.setattr(delegate_tool, "_get_child_timeout", lambda: timeout) + if hasattr(delegate_tool, "_get_worktree_isolation"): + monkeypatch.setattr(delegate_tool, "_get_worktree_isolation", lambda: False) + return delegate_tool._run_single_child( + task_index=0, + goal="exercise timeout transport drain", + child=child, + parent_agent=parent, + ) + + +def test_timeout_drains_transports_so_blocked_worker_can_unwind(monkeypatch): + child = _SslBlockedChild() + + result = _run(child, monkeypatch) + + assert result["status"] == "timeout" + # The drain ran from the timeout path (immediate sweep) and settled the + # blocked read; without it the worker would still be parked in ssl.read. + assert any(r.startswith("delegate_timeout") for r in child.drain_calls), ( + "timeout path never drained the abandoned child's transports" + ) + assert child.unwound.wait(timeout=5), ( + "worker never unwound — the drain did not settle its blocked read" + ) + assert child.closed.wait(timeout=5) + assert not child.close_while_blocked, ( + "child.close() ran while the worker was still inside its blocked read" + ) + + +def test_timeout_drain_failure_does_not_break_timeout_result(monkeypatch): + child = _SslBlockedChild() + + def _raising_drain(*, reason: str) -> int: + child.drain_calls.append(reason) + raise RuntimeError("transport sweep exploded") + + child._drain_transports_after_abandonment = _raising_drain + + result = _run(child, monkeypatch) + + assert result["status"] == "timeout" + assert child.drain_calls, "drain hook was never attempted" + # Unblock the worker manually so the deferred close can run. + child.read_settled.set() + assert child.unwound.wait(timeout=5) + assert child.closed.wait(timeout=5) + + +def test_timeout_without_drain_hook_still_defers_close(monkeypatch): + """Children lacking the hook (test doubles, third-party agents) keep the + plain deferred-close behavior.""" + child = _SslBlockedChild() + # Shadow the hook with a non-callable: the timeout path must skip it. + child.__dict__["_drain_transports_after_abandonment"] = None + + result = _run(child, monkeypatch) + + assert result["status"] == "timeout" + assert not child.closed.is_set(), ( + "close must stay deferred while the worker future is running" + ) + child.read_settled.set() + assert child.unwound.wait(timeout=5) + assert child.closed.wait(timeout=5) + assert not child.close_while_blocked + + +class _FakeSocket: + def __init__(self): + self.shutdown_calls = 0 + self.closed = False + + def settimeout(self, _v): + pass + + def shutdown(self, _how): + self.shutdown_calls += 1 + + def close(self): + self.closed = True + + +def test_agent_drain_shuts_sockets_down_without_fd_release(monkeypatch): + """AIAgent._drain_transports_after_abandonment must shutdown(), not close().""" + import threading as _threading + from unittest.mock import patch + + with patch("run_agent.AIAgent.__init__", return_value=None): + from run_agent import AIAgent + + agent = AIAgent.__new__(AIAgent) + + sock = _FakeSocket() + close_calls = {"n": 0} + + class _FakeClient: + def close(self): + close_calls["n"] += 1 + + agent.client = _FakeClient() + agent._client_lock = _threading.RLock() + agent._codex_session = None + agent._active_request_abort = None + + import agent.agent_runtime_helpers as arh + + monkeypatch.setattr(arh, "_iter_pool_sockets", lambda _c: iter([sock])) + + drained = agent._drain_transports_after_abandonment(reason="delegate_timeout_test") + + assert drained == 1 + assert sock.shutdown_calls == 1 + assert not sock.closed, "drain must never release socket FDs" + assert close_calls["n"] == 0, "drain must never call client.close()" diff --git a/tools/delegate_tool.py b/tools/delegate_tool.py index 770220e905..d295068843 100644 --- a/tools/delegate_tool.py +++ b/tools/delegate_tool.py @@ -3108,6 +3108,41 @@ def _run_single_child( _child_future.add_done_callback(_close_after_timed_out_worker) _child_close_deferred = True + + # Bounded drain (#94248 native half): the deferred close above + # only fires once the abandoned worker unwinds, but that worker + # is typically parked inside an in-flight OpenSSL read (Codex / + # httpx). Never hard-close that transport from this thread — + # releasing FDs under a live SSL read is the #29507/#70773 + # native-corruption family. Instead shutdown() the child's + # pooled sockets, which is FD-safe from any thread and settles + # the blocked read with EOF/EPIPE so the worker can unwind and + # trigger the deferred close. One immediate sweep plus one + # delayed re-sweep (covers a fresh connection opened between + # the interrupt and the first sweep); a worker that still + # doesn't settle keeps its resources until process exit rather + # than risking a cross-thread FD release. + _drain = getattr(child, "_drain_transports_after_abandonment", None) + if callable(_drain): + def _drain_once(phase: str) -> None: + try: + _drain(reason=f"delegate_timeout_{phase}") + except Exception: + logger.debug( + "Timed-out child transport drain (%s) failed", + phase, + exc_info=True, + ) + + _drain_once("immediate") + + def _drain_resweep() -> None: + if not _child_future.done(): + _drain_once("resweep") + + _resweep_timer = threading.Timer(5.0, _drain_resweep) + _resweep_timer.daemon = True + _resweep_timer.start() return _error_entry finally: # Shut down executor without waiting — if the child thread