fix(delegate): drain abandoned-worker transports FD-safely on child timeout
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
This commit is contained in:
69
run_agent.py
69
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.
|
||||
|
||||
|
||||
197
tests/tools/test_94248_timeout_transport_drain.py
Normal file
197
tests/tools/test_94248_timeout_transport_drain.py
Normal file
@@ -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()"
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user