diff --git a/gateway/run_agent_cache.py b/gateway/run_agent_cache.py index 2d1bfe5292..7fe1a26c53 100644 --- a/gateway/run_agent_cache.py +++ b/gateway/run_agent_cache.py @@ -383,19 +383,21 @@ class GatewayAgentCacheMixin: if interrupt_event is not None: interrupt_event._hermes_run_generation = int(generation) - async def _interrupt_and_clear_session( - self, session_key: str, source: SessionSource, *, interrupt_reason: str, - invalidation_reason: str, release_running_state: bool = True, - ) -> None: - """Interrupt the current run and clear queued session state consistently.""" + def _interrupt_running_turn(self, session_key: str, *, interrupt_reason: str, invalidation_reason: str) -> int: + """Sync core shared by /stop, /new and eviction: request a hard interrupt on the in-flight + agent, invalidate its run generation, and reap the tool processes that turn spawned. + Returns the post-bump generation.""" from gateway.run import _AGENT_PENDING_SENTINEL, _reap_gateway_turn_processes, request_hard_interrupt - if not session_key: - return state = self._peek_session_state(session_key) 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: - request_hard_interrupt(running_agent, interrupt_reason) + try: + 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: @@ -414,6 +416,19 @@ class GatewayAgentCacheMixin: name=f"gateway-turn-reaper-{_process_task_id[:12]}", daemon=True, ).start() + return _generation_at_interrupt + + async def _interrupt_and_clear_session( + self, session_key: str, source: SessionSource, *, interrupt_reason: str, + invalidation_reason: str, release_running_state: bool = True, + ) -> None: + """Interrupt the current run and clear queued session state consistently.""" + if not session_key: + return + state = self._peek_session_state(session_key) + self._interrupt_running_turn( + session_key, interrupt_reason=interrupt_reason, invalidation_reason=invalidation_reason, + ) adapter = self._adapter_for_source(source) interrupt_session_activity = getattr(type(adapter), "interrupt_session_activity", None) if adapter and callable(interrupt_session_activity): diff --git a/gateway/run_inbound.py b/gateway/run_inbound.py index 65ec84c22b..cd00061152 100644 --- a/gateway/run_inbound.py +++ b/gateway/run_inbound.py @@ -489,25 +489,11 @@ class GatewayInboundMixin: logger.debug("reaped-session staleness check failed", exc_info=True) def _hm_evict_running_agent(self, _quick_key: str, reason: str) -> None: - from gateway.run import _AGENT_PENDING_SENTINEL, _INTERRUPT_REASON_EVICTED, request_hard_interrupt - state = self._peek_session_state(_quick_key) - running_agent = state.turn.agent if state else None - if running_agent and running_agent is not _AGENT_PENDING_SENTINEL: - try: - request_hard_interrupt(running_agent, _INTERRUPT_REASON_EVICTED) - except Exception: - # Eviction must still invalidate and release the slot if a legacy or third-party - # agent cannot accept the interrupt request; otherwise the stale session remains - # unroutable and the pre-existing cleanup path is lost. - logger.warning( - "Failed to interrupt evicted agent for %s; continuing eviction cleanup", - _quick_key, - exc_info=True, - ) - self._invalidate_session_run_generation(_quick_key, reason=reason) + 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. + # releasing the slot so a late-finishing orphan cannot poison the replacement turn (#44212). self._evict_cached_agent(_quick_key) def _hm_merge_pending_for_source(