From 6e2796b82df8064a484bfcfe57b31c559752953c Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Thu, 3 Sep 2026 12:29:52 -0700 Subject: [PATCH] review-fix(locks): restore BASE fail-loud lock-slot reads and _is_openai_client_closed truth table MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Adolanium review §3 (Medium/Low). BASE shape at every steer/redirect/persist/agent-cache lock site was: lock = getattr(obj, "_x_lock", None); if lock is not None: with lock: ; else: . The refactor rewrote several of these as 'with getattr(...) or nullcontext(): getattr(slot, None)', which (a) turned a missing slot under the lock from a loud AttributeError into silent None and (b) in the gateway peek helper read the cache without any lock when the lock attribute was absent. Restored BASE semantics at: - agent/agent_runtime_helpers.py::_requeue_pending_steer - agent/interrupt_control.py: steer, redirect, clear_interrupt, _has_pending_redirect, _drain_pending_redirect, _drain_pending_steer (new _ic_slot helper: direct read under lock) - agent/session_persistence.py::_persist_lock (explicit None check + BASE rationale) - gateway/slash_commands.py::_cached_agent_for (BASE callers read ONLY under the lock; no lock -> None) - gateway/run_agent_cache.py::_evict_cached_agent (BASE: self._agent_cache direct under lock) - agent/client_lifecycle.py::_is_openai_client_closed: BASE body + docstring verbatim (outer is_closed first; inner _client.is_closed only when _client exists; else False) - agent/stream_delivery.py::_ensure_stream_writer_state: restore BASE rationale that the lock is created unconditionally in agent_init (_STREAM_STATE) and the lazy path is stub-only A/B (/tmp/rf/rev/ab_lock_fallbacks.py, ab_client_closed.py) is byte-identical BASE vs HEAD on the reviewer's inputs + edge cases. tests/agent/test_lock_fallback_base_semantics.py pins it (7 of 25 cases fail on the pre-fix tree). --- agent/agent_runtime_helpers.py | 12 +- agent/client_lifecycle.py | 28 ++- agent/interrupt_control.py | 21 +- agent/session_persistence.py | 9 +- agent/stream_delivery.py | 8 +- gateway/run_agent_cache.py | 13 +- gateway/slash_commands.py | 10 +- .../test_lock_fallback_base_semantics.py | 210 ++++++++++++++++++ 8 files changed, 288 insertions(+), 23 deletions(-) create mode 100644 tests/agent/test_lock_fallback_base_semantics.py diff --git a/agent/agent_runtime_helpers.py b/agent/agent_runtime_helpers.py index 870e46f479..e5d6e9c52a 100644 --- a/agent/agent_runtime_helpers.py +++ b/agent/agent_runtime_helpers.py @@ -3123,7 +3123,17 @@ def extract_api_error_context(error: Exception) -> Dict[str, Any]: def _requeue_pending_steer(agent, steer_text: str) -> None: """Put drained steer text back so the caller's fallback delivers it as a next-turn user message.""" - with getattr(agent, "_pending_steer_lock", None) or contextlib.nullcontext(): + # Under the lock the slot is read directly: an initialized agent always has both attributes, so a + # missing ``_pending_steer`` there is a real bug and must fail loud. The lock-less branch only + # exists for test stubs built via ``object.__new__`` that skipped ``__init__``. + _lock = getattr(agent, "_pending_steer_lock", None) + if _lock is not None: + with _lock: + if agent._pending_steer: + agent._pending_steer = agent._pending_steer + "\n" + steer_text + else: + agent._pending_steer = steer_text + else: existing = getattr(agent, "_pending_steer", None) agent._pending_steer = (existing + "\n" + steer_text) if existing else steer_text diff --git a/agent/client_lifecycle.py b/agent/client_lifecycle.py index fca0d90975..46b0d26448 100644 --- a/agent/client_lifecycle.py +++ b/agent/client_lifecycle.py @@ -79,13 +79,33 @@ class ClientLifecycleMixin: @staticmethod def _is_openai_client_closed(client: Any) -> bool: - """``is_closed`` is a property on httpx.Client but a method on openai.OpenAI (a bare getattr is always truthy).""" + """Check if an OpenAI client is closed. + + Handles both property and method forms of is_closed: + - httpx.Client.is_closed is a bool property + - openai.OpenAI.is_closed is a method returning bool + + Prior bug: getattr(client, "is_closed", False) returned the bound method, + which is always truthy, causing unnecessary client recreation on every call. + """ from unittest.mock import Mock + if isinstance(client, Mock): return False - is_closed = getattr(client, "is_closed", None) - closed = is_closed is not None and (is_closed() if callable(is_closed) else bool(is_closed)) - return bool(closed) or bool(getattr(getattr(client, "_client", None), "is_closed", False)) + + is_closed_attr = getattr(client, "is_closed", None) + if is_closed_attr is not None: + # Handle method (openai SDK) vs property (httpx) + if callable(is_closed_attr): + if is_closed_attr(): + return True + elif bool(is_closed_attr): + return True + + http_client = getattr(client, "_client", None) + if http_client is not None: + return bool(getattr(http_client, "is_closed", False)) + return False @staticmethod def _build_keepalive_http_client(base_url: str = "", *, verify: Any = True) -> Any: diff --git a/agent/interrupt_control.py b/agent/interrupt_control.py index e6ad163f44..e06f3ac13a 100644 --- a/agent/interrupt_control.py +++ b/agent/interrupt_control.py @@ -40,6 +40,15 @@ def _ic_lock(agent, attr: str): return contextlib.nullcontext() if lock is None else lock +def _ic_slot(agent, lock_attr: str, slot: str): + """Read the pending-text ``slot`` guarded by ``lock_attr``. An initialized agent always has both + attributes, so under the lock the slot is read directly and a missing one fails loud (a real bug); + only ``__init__``-less test stubs (no lock) get the ``getattr`` fallback.""" + if getattr(agent, lock_attr, None) is None: + return getattr(agent, slot, None) + return getattr(agent, slot) + + def _ic_codex_method(agent, name: str): """Codex app-server owns its model/tool loop; return its ``name`` hook or None.""" if getattr(agent, "api_mode", None) != "codex_app_server": @@ -188,7 +197,7 @@ class InterruptControlMixin: """Clear the interrupt request and per-thread tool signal. ``preserve_redirect`` is only for the conversation loop rebuilding the same logical turn after cancelling a model request.""" with _ic_lock(self, "_pending_redirect_lock"): - if preserve_redirect and not getattr(self, "_pending_redirect", None): + if preserve_redirect and not _ic_slot(self, "_pending_redirect_lock", "_pending_redirect"): return False self._interrupt_requested = False self._interrupt_message = self._tool_interrupt_reason = None @@ -211,7 +220,7 @@ class InterruptControlMixin: return False cleaned = text.strip() with _ic_lock(self, "_pending_steer_lock"): - existing = getattr(self, "_pending_steer", None) + existing = _ic_slot(self, "_pending_steer_lock", "_pending_steer") self._pending_steer = (existing + "\n" + cleaned) if existing else cleaned return True @@ -243,7 +252,7 @@ class InterruptControlMixin: with _ic_lock(self, "_pending_redirect_lock"): if _model_active is None or not _model_active.is_set(): return False # response completed before we got the lock: surface queues a new turn - existing = getattr(self, "_pending_redirect", None) + existing = _ic_slot(self, "_pending_redirect_lock", "_pending_redirect") if self._interrupt_requested and not existing: return False self._pending_redirect = ( @@ -265,18 +274,18 @@ class InterruptControlMixin: def _has_pending_redirect(self) -> bool: """Return whether an active-turn redirect is waiting to be applied.""" with _ic_lock(self, "_pending_redirect_lock"): - return bool(getattr(self, "_pending_redirect", None)) + return bool(_ic_slot(self, "_pending_redirect_lock", "_pending_redirect")) def _drain_pending_redirect(self) -> Optional[str]: """Return and clear pending active-turn correction text.""" with _ic_lock(self, "_pending_redirect_lock"): - text = getattr(self, "_pending_redirect", None) + text = _ic_slot(self, "_pending_redirect_lock", "_pending_redirect") self._pending_redirect = None return text def _drain_pending_steer(self) -> Optional[str]: """Return the pending steer text (if any) and clear the slot; None when nothing is pending.""" with _ic_lock(self, "_pending_steer_lock"): - text = getattr(self, "_pending_steer", None) + text = _ic_slot(self, "_pending_steer_lock", "_pending_steer") self._pending_steer = None return text diff --git a/agent/session_persistence.py b/agent/session_persistence.py index 3a392e4a55..773406520a 100644 --- a/agent/session_persistence.py +++ b/agent/session_persistence.py @@ -104,8 +104,13 @@ def _durable_content(content: Any) -> Any: def _persist_lock(agent): - """Close and turn-start persistence can run on separate CLI threads: one critical section.""" - return getattr(agent, "_session_persist_lock", None) or nullcontext() + """Close and turn-start persistence can run on separate CLI threads: one critical section. + + ``__init__`` always creates ``_session_persist_lock``; only ``object.__new__``-built test stubs lack it + (they run unlocked, matching the historical ``if persist_lock is None`` branch). + """ + lock = getattr(agent, "_session_persist_lock", None) + return nullcontext() if lock is None else lock # --- flush phases (module-level so the flush also works bound onto duck-typed agents) --- diff --git a/agent/stream_delivery.py b/agent/stream_delivery.py index 34b040df29..2f1cfb6ccb 100644 --- a/agent/stream_delivery.py +++ b/agent/stream_delivery.py @@ -203,9 +203,13 @@ class StreamDeliveryMixin: self._deliver_interim(visible, already_streamed=already_streamed, record=undelivered_parts or [visible]) def _ensure_stream_writer_state(self) -> None: - """Lazily create the single-writer guard fields (``AIAgent.__new__``-built instances skip ``agent_init``). + """Lazily create the single-writer guard fields (#65991). - See #65991. + The fields are normally set unconditionally in ``agent_init`` (``_STREAM_STATE``), so every + ``__init__``-built agent shares ONE lock from birth and the lazy path below is never taken by + two threads. Only agents constructed via ``AIAgent.__new__`` (test doubles, legacy/partially- + initialized instances) skip that path; claiming/checking the writer must not crash those, so + initialize the fields on first use. """ if getattr(self, "_stream_writer_lock", None) is None: self._stream_writer_lock = threading.Lock() diff --git a/gateway/run_agent_cache.py b/gateway/run_agent_cache.py index 717076129d..fdf701a7cf 100644 --- a/gateway/run_agent_cache.py +++ b/gateway/run_agent_cache.py @@ -581,11 +581,16 @@ class GatewayAgentCacheMixin: if state is not None: state.conversation.ephemeral_pin = None state.conversation.vc_last = None - # Tests build runners with ``_agent_cache_lock = None``; evict lock-free then. - _cache = getattr(self, "_agent_cache", None) + # Tests build runners with ``_agent_cache_lock = None``; evict lock-free then. With the lock + # present ``_agent_cache`` is read directly (an initialized runner always has it). + _lock = getattr(self, "_agent_cache_lock", None) evicted = None - if _cache is not None: - with getattr(self, "_agent_cache_lock", None) or nullcontext(): + if _lock: + with _lock: + evicted = self._agent_cache.pop(session_key, None) + else: + _cache = getattr(self, "_agent_cache", None) + if _cache is not None: evicted = _cache.pop(session_key, None) agent = _first_agent(evicted) # Never tear down an agent that's mid-turn — its client, sandbox and child subagents are in use. diff --git a/gateway/slash_commands.py b/gateway/slash_commands.py index 547c471957..76beb5270a 100644 --- a/gateway/slash_commands.py +++ b/gateway/slash_commands.py @@ -169,13 +169,15 @@ class GatewaySlashCommandsMixin( # ------------------------------------------------------------------ shared helpers def _cached_agent_for(self, session_key: str): """Peek the cached AIAgent for *session_key* without evicting it, or None. Entries are - ``(agent, signature, ...)`` tuples (bare agents from test doubles accepted); lock/cache may - be absent on fixtures that skip ``__init__``.""" + ``(agent, signature, ...)`` tuples (bare agents from test doubles accepted). Every historical + caller read the cache ONLY under ``_agent_cache_lock``; fixtures that skip ``__init__`` (no lock + or no cache) get None rather than an unlocked read.""" cache = getattr(self, "_agent_cache", None) - if cache is None: + lock = getattr(self, "_agent_cache_lock", None) + if cache is None or lock is None: return None try: - with getattr(self, "_agent_cache_lock", None) or contextlib.nullcontext(): + with lock: entry = cache.get(session_key) except Exception: return None diff --git a/tests/agent/test_lock_fallback_base_semantics.py b/tests/agent/test_lock_fallback_base_semantics.py new file mode 100644 index 0000000000..eca7dc6cc5 --- /dev/null +++ b/tests/agent/test_lock_fallback_base_semantics.py @@ -0,0 +1,210 @@ +"""Pin BASE (main) semantics at the lock-fallback sites and the ``_is_openai_client_closed`` truth table. + +Review-fix for the simplify-codebase PR: the refactor briefly rewrote several ``lock = getattr(agent, +"_x_lock", None); if lock is not None: with lock: `` blocks as ``with getattr(...) or +nullcontext(): getattr(agent, slot, None)``, which (a) silently runs the critical section unlocked and (b) +turns a missing slot under the lock from a loud AttributeError into a silent None. These tests hold the +original shape: an initialized agent (lock present) reads its slots directly and fails loud; only +``object.__new__`` stubs (no lock) get the getattr fallback. +""" +from __future__ import annotations + +import threading +from unittest.mock import Mock + +import pytest + +from agent.agent_runtime_helpers import _requeue_pending_steer +from agent.client_lifecycle import ClientLifecycleMixin +from agent.interrupt_control import InterruptControlMixin +from agent.session_persistence import SessionPersistenceMixin +from agent.stream_delivery import StreamDeliveryMixin +from gateway.run_agent_cache import GatewayAgentCacheMixin +from gateway.slash_commands import GatewaySlashCommandsMixin + + +class _RecLock: + def __init__(self): + self.entered = 0 + self._lock = threading.Lock() + + def __enter__(self): + self.entered += 1 + return self._lock.__enter__() + + def __exit__(self, *exc): + return self._lock.__exit__(*exc) + + +class _Bare: + pass + + +# ---------------------------------------------------------------- steer requeue + + +def test_requeue_pending_steer_locked_path_reads_slot_directly_and_fails_loud(): + agent = _Bare() + agent._pending_steer_lock = _RecLock() + with pytest.raises(AttributeError): + _requeue_pending_steer(agent, "new") # lock present but slot missing => real bug, loud + agent._pending_steer = "old" + _requeue_pending_steer(agent, "new") + assert agent._pending_steer == "old\nnew" + assert agent._pending_steer_lock.entered >= 1 + + +def test_requeue_pending_steer_unlocked_stub_fallback(): + agent = _Bare() # object.__new__-style stub: no lock, no slot + _requeue_pending_steer(agent, "new") + assert agent._pending_steer == "new" + + +# ------------------------------------------------------- interrupt_control slots + + +@pytest.mark.parametrize( + "method, lock_attr, slot", + [ + ("steer", "_pending_steer_lock", "_pending_steer"), + ("_drain_pending_steer", "_pending_steer_lock", "_pending_steer"), + ("_has_pending_redirect", "_pending_redirect_lock", "_pending_redirect"), + ("_drain_pending_redirect", "_pending_redirect_lock", "_pending_redirect"), + ], +) +def test_interrupt_control_slot_reads_fail_loud_only_under_lock(method, lock_attr, slot): + fn = getattr(InterruptControlMixin, method) + call = (lambda a: fn(a, "new")) if method == "steer" else fn + + locked = _Bare() + setattr(locked, lock_attr, threading.Lock()) + with pytest.raises(AttributeError): + call(locked) + + unlocked = _Bare() + call(unlocked) # no lock: stub fallback, no raise + + setattr(locked, slot, "old") + call(locked) # slot present: normal path + + +def test_steer_and_drain_roundtrip_with_lock(): + agent = _Bare() + agent._pending_steer_lock = threading.Lock() + agent._pending_steer = None + assert InterruptControlMixin.steer(agent, "a") + assert InterruptControlMixin.steer(agent, "b") + assert InterruptControlMixin._drain_pending_steer(agent) == "a\nb" + assert agent._pending_steer is None + + +# ---------------------------------------------------------- session persistence + + +def test_flush_messages_uses_lock_when_present_and_runs_unlocked_for_stubs(): + class P(_Bare): + def _flush_messages_to_session_db_unlocked(self, messages, conversation_history=None): + return ("called", len(messages)) + + locked = P() + locked._session_persist_lock = _RecLock() + assert SessionPersistenceMixin._flush_messages_to_session_db(locked, [{"role": "user"}]) == ("called", 1) + assert locked._session_persist_lock.entered == 1 + + stub = P() + assert SessionPersistenceMixin._flush_messages_to_session_db(stub, [{"role": "user"}]) == ("called", 1) + + +# ---------------------------------------------------------------- stream writer + + +def test_stream_writer_lock_is_created_unconditionally_at_init(): + from agent.agent_init import _STREAM_STATE + + assert _STREAM_STATE["_stream_writer_lock"] is threading.Lock + # The lazy path is only for __new__-built stubs and never replaces an existing lock. + agent = StreamDeliveryMixin.__new__(StreamDeliveryMixin) + agent._stream_writer_lock = existing = threading.Lock() + StreamDeliveryMixin._ensure_stream_writer_state(agent) + assert agent._stream_writer_lock is existing + + +# ------------------------------------------------------- gateway agent-cache sites + + +class _Runner(_Bare): + def _peek_session_state(self, key): + return None + + def _running_agent_ids(self): + return set() + + def _spawn_release_thread(self, *a, **k): + pass + + def _release_evicted_agent_soft(self, *a): + pass + + +def test_cached_agent_for_reads_only_under_lock(): + runner = _Runner() + runner._agent_cache = {"k": ("AGENT", "sig")} + assert GatewaySlashCommandsMixin._cached_agent_for(runner, "k") is None # no lock: no unlocked read + runner._agent_cache_lock = _RecLock() + assert GatewaySlashCommandsMixin._cached_agent_for(runner, "k") == "AGENT" + assert runner._agent_cache_lock.entered == 1 + + +def test_evict_cached_agent_reads_cache_directly_when_locked(): + runner = _Runner() + runner._agent_cache_lock = threading.Lock() + with pytest.raises(AttributeError): + GatewayAgentCacheMixin._evict_cached_agent(runner, "k") # lock present, cache missing => loud + runner._agent_cache = {"k": ("AGENT", "sig")} + GatewayAgentCacheMixin._evict_cached_agent(runner, "k") + assert runner._agent_cache == {} + stub = _Runner() + stub._agent_cache = {"k": ("AGENT", "sig")} + GatewayAgentCacheMixin._evict_cached_agent(stub, "k") # lock-less test runner: evicts lock-free + assert stub._agent_cache == {} + + +# ------------------------------------------------------ _is_openai_client_closed + + +class _Inner: + def __init__(self, is_closed): + self.is_closed = is_closed + + +def _duck(outer=None, inner=None, *, has_outer=True, has_inner=True): + d = _Bare() + if has_outer: + d.is_closed = outer + if has_inner: + d._client = inner + return d + + +@pytest.mark.parametrize( + "client, expected", + [ + # the reviewer's 4 combos: duck outer is_closed True/False, with/without _client + (_duck(True, has_inner=False), True), + (_duck(False, has_inner=False), False), + (_duck(True, _Inner(False)), True), # outer says closed => closed, inner not consulted + (_duck(False, _Inner(True)), True), # outer open, inner closed => inner is the answer + (_duck(False, _Inner(False)), False), + (_duck(lambda: True, _Inner(False)), True), # openai.OpenAI.is_closed() method form + (_duck(lambda: False, _Inner(True)), True), + (_duck(lambda: False, _Inner(False)), False), + (_duck(lambda: False, has_inner=False), False), + (_duck(None, _Inner(True)), True), # outer None => fall through to inner + (_duck(None, None), False), + (_duck(has_outer=False, inner=_Inner(True)), True), + (_duck(has_outer=False, has_inner=False), False), + (Mock(), False), + ], +) +def test_is_openai_client_closed_truth_table(client, expected): + assert ClientLifecycleMixin._is_openai_client_closed(client) is expected