refactor(gateway): drop the write-only one-shot generation stamp; /stop and eviction share _drop_turn_slot

- The `one_turn_restore["run_generation"]` stamp had no reader once settlement
  moved to the invalidate chokepoint (the finalizer guards on its own generation
  via _is_session_run_current); delete it and the test lines that set it.
- release-slot-then-evict-cached-agent (#44212 rationale) was duplicated in
  /stop and eviction; one _drop_turn_slot owns it.
- Best-effort interrupt uses the repo's _log_suppressed seam like run_shutdown.
- turn_lease.rebind resolves the lease via token.lease like release does
  (identity, not a session_id lookup); _held_turn_lease hands back the token
  map so release/rebind stop re-peeking session state.
- Test trims: vacuous isinstance, registry-internals asserts.
This commit is contained in:
kshitijk4poor
2026-09-11 12:46:43 +05:30
committed by kshitij
parent 90970d85f2
commit 789eec791b
6 changed files with 30 additions and 42 deletions

View File

@@ -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

View File

@@ -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

View File

@@ -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:

View File

@@ -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

View File

@@ -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:

View File

@@ -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