review-fix(locks): restore BASE fail-loud lock-slot reads and _is_openai_client_closed truth table

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: <direct
attribute read>; else: <getattr fallback for object.__new__ test stubs>. 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).
This commit is contained in:
Teknium
2026-09-03 12:29:52 -07:00
parent 283c4f058c
commit 6e2796b82d
8 changed files with 288 additions and 23 deletions
+11 -1
View File
@@ -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
+24 -4
View File
@@ -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:
+15 -6
View File
@@ -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
+7 -2
View File
@@ -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) ---
+6 -2
View File
@@ -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()
+9 -4
View File
@@ -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.
+6 -4
View File
@@ -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
@@ -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: <direct attr read>`` 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