diff --git a/gateway/run.py b/gateway/run.py index aa74b820e1..61d795d64f 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -3518,11 +3518,8 @@ class GatewayRunner( # to one session_id (switch_session's many-to-one mapping), which routing-key guards cannot see. self._turn_leases = SessionTurnLeaseRegistry() # Stall-notified keys clear when pending clears / activity resumes / conversation boundary. - # Tokens for held turn leases, keyed by (routing key, run generation) so release is granted per-turn - # and a stale unwind can never free a newer turn's lease (#28686 ownership lesson). Held turn-lease - # tokens live on SessionState.turn.lease_tokens keyed by run generation (the old dict was keyed - # (routing key, generation) so a stale unwind could never free a newer turn's lease — the - # per-generation key preserves that ownership check, #28686). Runner-level queued interrupt text lives on + # Held turn-lease tokens live on SessionState.turn.lease_tokens keyed by run generation, so a + # stale unwind can never free a newer turn's lease (#28686). Runner-level queued interrupt text lives on # SessionState.persistent.pending_command_text (NOTE: distinct from the adapter-level # _pending_messages Dict[str, MessageEvent] in gateway/platforms/base.py, which shares the legacy # name). Last successfully-resolved (non-empty) model, keyed by session. Used as a fallback when a diff --git a/gateway/run_agent_cache.py b/gateway/run_agent_cache.py index 7fe1a26c53..f2ddf9e64b 100644 --- a/gateway/run_agent_cache.py +++ b/gateway/run_agent_cache.py @@ -14,6 +14,7 @@ from typing import TYPE_CHECKING, Any, Dict, List, Optional from agent.interrupt_compat import _accepts_keyword from gateway.config import Platform from gateway.session import SessionSource, build_session_context_prompt +from gateway.run_shutdown import _log_suppressed from hermes_cli.config import cfg_get if TYPE_CHECKING: # string annotations only; never imported at runtime (cycle) @@ -255,14 +256,24 @@ class GatewayAgentCacheMixin: self._persist_active_agents() return True + def _drop_turn_slot(self, session_key: str) -> None: + """Release the running-agent slot and evict the cached instance (/stop, eviction, reaper). + ``_interrupt_requested`` is cleared only by the turn finalizer, so on a hung/still-draining + run the flag would survive and silently kill the session's NEXT message (interrupted=True, + api_calls=0, empty response); the next message rebuilds from history while the old agent + keeps its flag so a hung drain still dies (#44212).""" + self._release_running_agent_state(session_key) + self._evict_cached_agent(session_key) + def _held_turn_lease(self, session_key: str, run_generation: int): - """Return ``(registry, token)`` when ``session_key`` holds a lease token for ``run_generation``, else None.""" + """Return ``(registry, lease_tokens)`` when ``session_key`` holds a lease token for + ``run_generation``, else None. Callers ``get``/``pop`` the token from the map themselves.""" registry = getattr(self, "_turn_leases", None) state = self._peek_session_state(session_key) if session_key and registry is not None else None - token = state.turn.lease_tokens.get(run_generation) if state is not None else None - if token is None: + tokens = state.turn.lease_tokens if state is not None else None + if tokens is None or run_generation not in tokens: return None - return registry, token + return registry, tokens def _release_turn_lease(self, session_key: str, run_generation: int) -> bool: """Release the turn lease acquired by (``session_key``, ``run_generation``). Keyed by (routing @@ -271,8 +282,8 @@ class GatewayAgentCacheMixin: held = self._held_turn_lease(session_key, run_generation) if held is None: return False - registry, token = held - self._peek_session_state(session_key).turn.lease_tokens.pop(run_generation, None) + registry, tokens = held + token = tokens.pop(run_generation) try: return registry.release(token) except Exception: @@ -286,9 +297,9 @@ class GatewayAgentCacheMixin: held = self._held_turn_lease(session_key, run_generation) if new_session_id else None if held is None: return False - registry, token = held + registry, tokens = held try: - return registry.rebind(token, new_session_id) + return registry.rebind(tokens[run_generation], new_session_id) except Exception: logger.debug("Failed to rebind turn lease", exc_info=True) return False @@ -392,12 +403,11 @@ class GatewayAgentCacheMixin: running_agent = state.turn.agent if state else None _process_task_id, _process_baseline = "", None if running_agent and running_agent is not _AGENT_PENDING_SENTINEL: - try: + # A raising interrupt implementation must not leave the slot unroutable: the generation + # bump and release below are the cleanup that matters. + with _log_suppressed(logging.WARNING, "Failed to interrupt running agent for %s; continuing", + session_key, exc_info=True): request_hard_interrupt(running_agent, interrupt_reason) - except Exception: - # A raising interrupt implementation must not leave the slot unroutable: the - # generation bump and release below are the cleanup that matters. - logger.warning("Failed to interrupt running agent for %s; continuing", session_key, exc_info=True) _process_task_id = getattr(running_agent, "_gateway_turn_process_task_id", "") _process_baseline = getattr(running_agent, "_gateway_turn_process_baseline", None) # Bump the generation BEFORE scheduling the reap thread and capture the post-bump value: @@ -442,13 +452,7 @@ class GatewayAgentCacheMixin: if state is not None: state.persistent.pending_command_text = None if release_running_state: - self._release_running_agent_state(session_key) - # Evict the cached agent: ``_interrupt_requested`` is only cleared by the turn finalizer, - # so on a hung/still-draining run the flag survives and silently kills the session's NEXT - # message (interrupted=True, api_calls=0, empty response). Like /new and /model, the next - # message rebuilds from history; the old agent keeps its flag so a hung drain still dies. - # See #44212. - self._evict_cached_agent(session_key) + self._drop_turn_slot(session_key) async def _refresh_agent_cache_message_count(self, session_key: str, session_id: Optional[str]) -> None: """Re-baseline a cached agent's stored message_count after THIS turn — the coherence guard diff --git a/gateway/run_inbound.py b/gateway/run_inbound.py index cd00061152..7dcd4cdf91 100644 --- a/gateway/run_inbound.py +++ b/gateway/run_inbound.py @@ -491,10 +491,7 @@ class GatewayInboundMixin: def _hm_evict_running_agent(self, _quick_key: str, reason: str) -> None: from gateway.run import _INTERRUPT_REASON_EVICTED self._interrupt_running_turn(_quick_key, interrupt_reason=_INTERRUPT_REASON_EVICTED, invalidation_reason=reason) - self._release_running_agent_state(_quick_key) - # The interrupt flag is cleared only by the turn finalizer. Remove the cached instance after - # releasing the slot so a late-finishing orphan cannot poison the replacement turn (#44212). - self._evict_cached_agent(_quick_key) + self._drop_turn_slot(_quick_key) def _hm_merge_pending_for_source( self, source: SessionSource, _quick_key: str, event: "MessageEvent", *, merge_text: bool = False @@ -1263,11 +1260,6 @@ class GatewayInboundMixin: _claim_state.turn.started_ts = time.time() self._persist_active_agents() _run_generation = self._begin_session_run_generation(_quick_key) - # A pending one-shot snapshot (/moa, /model --once) belongs to the turn that claims the slot: - # only its finalizer may restore it. Stop/reset/eviction settle it before bumping the - # generation (``_invalidate_session_run_generation``), so a displaced turn never leaks it. - if _claim_state.conversation.one_turn_restore: - _claim_state.conversation.one_turn_restore["run_generation"] = _run_generation try: try: diff --git a/gateway/turn_lease.py b/gateway/turn_lease.py index e04355dba0..6504a8ed7a 100644 --- a/gateway/turn_lease.py +++ b/gateway/turn_lease.py @@ -144,7 +144,8 @@ class SessionTurnLeaseRegistry: if (token is None or token.released or not new_session_id or new_session_id == token.session_id): return False - if (lease := self._leases.get(token.session_id)) is None or lease.holder is not token: + lease = token.lease + if lease.holder is not token: return False existing = self._leases.get(new_session_id) if existing is not None and existing is not lease and not existing.idle: @@ -156,7 +157,6 @@ class SessionTurnLeaseRegistry: token.session_id, new_session_id, token.owner_key, token.generation, *_holder_desc(existing.holder), new_session_id) return False - # Preserve alias across rotation: both old and new session IDs point to the same lease self._leases[new_session_id] = lease lease.last_used = time.time() token.session_id = new_session_id diff --git a/tests/gateway/test_moa_one_shot_restore.py b/tests/gateway/test_moa_one_shot_restore.py index 1f36cd9338..e9a32df633 100644 --- a/tests/gateway/test_moa_one_shot_restore.py +++ b/tests/gateway/test_moa_one_shot_restore.py @@ -27,7 +27,6 @@ def _runner_with_pending_once(): def test_restore_runs_from_finally_even_when_turn_raises(): runner, state = _runner_with_pending_once() gen = runner._begin_session_run_generation(KEY) - state.conversation.one_turn_restore["run_generation"] = gen with pytest.raises(RuntimeError): try: diff --git a/tests/gateway/test_reaped_eviction_interrupts_run.py b/tests/gateway/test_reaped_eviction_interrupts_run.py index c0f052a418..76e407ef46 100644 --- a/tests/gateway/test_reaped_eviction_interrupts_run.py +++ b/tests/gateway/test_reaped_eviction_interrupts_run.py @@ -77,13 +77,13 @@ def test_eviction_interrupts_before_release_and_drops_cached_agent(entrypoint: s else: gateway._hm_evict_running_agent(KEY, "stale_running_agent_eviction") - assert isinstance(gateway, GatewayInboundMixin) assert agent.interrupted assert events[0] == ("interrupt", _INTERRUPT_REASON_EVICTED, True) release_events = [event for event in events if event[0] == "release"] assert release_events == [("release", True)] # the interrupt was requested BEFORE the slot release assert gateway._peek_session_state(KEY).turn.agent is None assert KEY not in gateway._agent_cache + # The reason must be a registered control message or the finalizer treats it as user text. assert _is_control_interrupt_message(_INTERRUPT_REASON_EVICTED) @@ -127,7 +127,6 @@ def test_one_shot_override_settles_on_stop_and_stale_finalizer_is_a_noop() -> No state.conversation.model_override = {"model": "once-model", "provider": "test"} state.conversation.one_turn_restore = {"had_override": True, "override": dict(prior)} owning_gen = gateway._begin_session_run_generation(KEY) - state.conversation.one_turn_restore["run_generation"] = owning_gen gateway._invalidate_session_run_generation(KEY, reason="user_stop") # settlement point assert state.conversation.model_override == prior @@ -152,9 +151,6 @@ async def test_turn_lease_rebind_preserves_parent_lock_domain_and_releases() -> assert registry.rebind(token, "child-session") is True assert token.session_id == "child-session" - # Both parent and child session IDs must be registered to the same lease - assert registry._leases.get("parent-session") is registry._leases.get("child-session") - assert registry._leases["parent-session"].holder is token # Parent lock domain remains busy while child is held import asyncio