diff --git a/gateway/run_shutdown.py b/gateway/run_shutdown.py index 52bf07449a..dcd93ff693 100644 --- a/gateway/run_shutdown.py +++ b/gateway/run_shutdown.py @@ -1055,16 +1055,17 @@ class GatewayShutdownMixin: flush_agent_history_to_file(getattr(agent, "session_id", None), _session_messages) async def _finalize_shutdown_agents(self, active_agents: Dict[str, Any]) -> None: - for agent in active_agents.values(): + for session_key, agent in active_agents.items(): self._flush_agent_transcript_at_shutdown(agent) # Off-loop + bounded: plugin on_session_finalize hooks can do arbitrary synchronous work # (e.g. a full-session trace export) — same hang class as the memory provider below. await self._finalize_session_off_loop( session_id=getattr(agent, "session_id", None), platform="gateway", reason="shutdown", + session_key=session_key, ) # Off-loop + bounded: a wedged memory provider here used to hang the whole shutdown so # SIGTERM never completed. - await self._cleanup_agent_resources_off_loop(agent, context="shutdown finalize") + await self._cleanup_agent_resources_off_loop(agent, context="shutdown finalize", session_key=session_key) def _should_emit_long_running_notification( self, session_key: Optional[str], agent: Any, executor_task: Optional[Any], @@ -1109,15 +1110,23 @@ class GatewayShutdownMixin: tasks = self._deferred_agent_cleanup_tasks = set() self._track_task_in(tasks, asyncio.create_task(_cleanup_when_done())) - async def _finalize_session_off_loop(self, *, session_id: Any, platform: str, reason: str, **extra: Any) -> None: - """Run hermes_cli.lifecycle.finalize_session off-loop, bounded; on timeout the worker is left alone.""" + async def _finalize_session_off_loop( + self, *, session_id: Any, platform: str, reason: str, session_key: Optional[str] = None, **extra: Any, + ) -> None: + """Run hermes_cli.lifecycle.finalize_session off-loop, bounded; on timeout the worker is left alone. + ``session_key`` lets an unscoped caller (shutdown) enter the owning profile's scope: plugin + ``on_session_finalize`` observers and the Relay coordinator (``current_profile_key``) resolve + profile state at call time.""" def _call() -> None: from hermes_cli.lifecycle import finalize_session finalize_session(session_id=session_id, platform=platform, reason=reason, **extra) try: - await asyncio.wait_for(self._run_in_executor_with_context(_call), timeout=self._FINALIZE_TIMEOUT_S) + await asyncio.wait_for( + self._run_in_executor_with_context(self._run_release_in_profile_scope, _call, (), session_key), + timeout=self._FINALIZE_TIMEOUT_S, + ) except asyncio.TimeoutError: logger.warning( "Session finalize hooks (%s, reason=%s) exceeded %ss; proceeding without blocking the event loop " @@ -1126,8 +1135,17 @@ class GatewayShutdownMixin: except Exception as finalize_exc: logger.debug("Session finalize hooks (%s, reason=%s) failed: %s", session_id, reason, finalize_exc) - async def _cleanup_agent_resources_off_loop(self, agent: Any, *, context: str = "") -> None: - """Run _cleanup_agent_resources in a worker thread, bounded; on timeout the worker is left alone.""" + async def _cleanup_agent_resources_off_loop( + self, agent: Any, *, context: str = "", session_key: Optional[str] = None, + ) -> None: + """Run _cleanup_agent_resources in a worker thread, bounded; on timeout the worker is left alone. + + The teardown fires the memory-provider lifecycle hooks (``flush_pending`` → ``on_session_end`` → + ``shutdown`` → ``close``), which read credentials/home at call time. In-turn callers carry the + profile scope through ``_run_in_executor_with_context``; shutdown does not (it runs on the main + loop, outside any adapter handler), so under multiplexing ``on_session_end`` failed closed and the + session tail was never committed (#110622). ``_run_release_in_profile_scope`` enters the OWNING + profile's scope from ``session_key`` when the caller has none, exactly like cache eviction.""" if agent is None: return if context.startswith("shutdown") or context == "session expiry": @@ -1136,7 +1154,9 @@ class GatewayShutdownMixin: ctx_label = f" ({context})" if context else "" try: await asyncio.wait_for( - self._run_in_executor_with_context(self._cleanup_agent_resources, agent), + self._run_in_executor_with_context( + self._run_release_in_profile_scope, self._cleanup_agent_resources, (agent,), session_key, + ), timeout=self._CLEANUP_TIMEOUT_S, ) except asyncio.TimeoutError: @@ -1717,12 +1737,13 @@ class GatewayShutdownMixin: _cache = getattr(self, "_agent_cache", None) if _cache_lock is not None and _cache is not None: with _cache_lock: - _idle_agents = list(_cache.values()) + _idle_agents = list(_cache.items()) _cache.clear() - for _entry in _idle_agents: + for _key, _entry in _idle_agents: # Bounded + off-loop: a wedged memory provider here once made SIGTERM hang forever. await self._cleanup_agent_resources_off_loop( - _entry[0] if isinstance(_entry, tuple) else _entry, context="shutdown idle-cache" + _entry[0] if isinstance(_entry, tuple) else _entry, context="shutdown idle-cache", + session_key=_key, ) # Settle completion flush tasks while adapters are alive so every watcher gets a retryable result. cancel_completion_batches = getattr(self, "_cancel_process_completion_batch_tasks", None) diff --git a/tests/gateway/test_finalize_session_off_loop.py b/tests/gateway/test_finalize_session_off_loop.py index cf10d13617..e209ebae9c 100644 --- a/tests/gateway/test_finalize_session_off_loop.py +++ b/tests/gateway/test_finalize_session_off_loop.py @@ -133,7 +133,7 @@ def test_shutdown_finalize_path_uses_off_loop_dispatch(monkeypatch): GatewayRunner, "_finalize_session_off_loop", _fake_off_loop ) - async def _fake_cleanup(self, agent, *, context=""): + async def _fake_cleanup(self, agent, *, context="", session_key=None): return None monkeypatch.setattr( diff --git a/tests/gateway/test_shutdown_cache_cleanup.py b/tests/gateway/test_shutdown_cache_cleanup.py index 71b77e47a4..3123bb825f 100644 --- a/tests/gateway/test_shutdown_cache_cleanup.py +++ b/tests/gateway/test_shutdown_cache_cleanup.py @@ -76,7 +76,7 @@ class _FakeGateway: # inline in tests so the bounded-cleanup path is exercised. return func(*args) - async def _cleanup_agent_resources_off_loop(self, agent, *, context=""): + async def _cleanup_agent_resources_off_loop(self, agent, *, context="", session_key=None): # Mirror the real bounded helper, inline (no executor/timeout) so the # fake exercises the same call shape stop() now uses. self._cleanup_agent_resources(agent) diff --git a/tests/gateway/test_shutdown_teardown_profile_scope.py b/tests/gateway/test_shutdown_teardown_profile_scope.py new file mode 100644 index 0000000000..c44a4eb76d --- /dev/null +++ b/tests/gateway/test_shutdown_teardown_profile_scope.py @@ -0,0 +1,88 @@ +"""Multiplex invariant: gateway-hosted teardown fires memory-provider lifecycle hooks under the +OWNING profile's scope (#110622). + +``on_session_end`` / ``shutdown`` / ``close`` read credentials and home at call time. Shutdown +runs on the main loop outside any adapter handler, so its executor hop carried an EMPTY scope and +a secondary profile's ``get_secret("OPENVIKING_API_KEY")`` failed closed — the end-of-session +commit was skipped and logged as a WARNING + traceback on every gateway-hosted session close. +""" +from __future__ import annotations + +import asyncio +from contextvars import copy_context +from pathlib import Path +from types import SimpleNamespace + +import pytest + +from agent import secret_scope +from gateway.config import GatewayConfig +from gateway.run import GatewayRunner +from hermes_constants import get_hermes_home + + +def _runner(profile_homes: dict[str, Path]) -> GatewayRunner: + runner = object.__new__(GatewayRunner) + runner.config = GatewayConfig(multiplex_profiles=True) + runner.session_store = SimpleNamespace(_profile_home_for_key=lambda key: profile_homes.get(key)) + from concurrent.futures import ThreadPoolExecutor + executor = ThreadPoolExecutor(max_workers=2) + + async def _run_in_executor_with_context(func, *args): + loop = asyncio.get_running_loop() + ctx = copy_context() + return await loop.run_in_executor(executor, lambda: ctx.run(func, *args)) + + runner._run_in_executor_with_context = _run_in_executor_with_context + return runner + + +def _recording_agent(seen: dict): + class _Provider: + name = "recorder" + + def on_session_end(self, messages): + seen["scope"] = secret_scope.current_secret_scope() + seen["home"] = get_hermes_home() + + def shutdown_memory_provider(messages=None): + _Provider().on_session_end(messages) + + return SimpleNamespace( + shutdown_memory_provider=shutdown_memory_provider, close=lambda: None, _session_messages=[], + ) + + +@pytest.mark.asyncio +async def test_shutdown_teardown_commits_memory_under_the_owning_profile(tmp_path, monkeypatch): + default_home = tmp_path / ".hermes" + prof_b = default_home / "profiles" / "b" + prof_b.mkdir(parents=True) + (prof_b / ".env").write_text("OPENVIKING_API_KEY=key-of-b\n") + monkeypatch.setenv("HERMES_HOME", str(default_home)) + seen: dict = {} + secret_scope.set_multiplex_active(True) + try: + # Shutdown runs on the main loop with NO scope installed. + assert secret_scope.current_secret_scope() is None + await _runner({"agent:b:telegram:dm:1": prof_b})._cleanup_agent_resources_off_loop( + _recording_agent(seen), context="shutdown finalize", session_key="agent:b:telegram:dm:1", + ) + finally: + secret_scope.set_multiplex_active(False) + assert seen["home"] == prof_b + assert seen["scope"] and seen["scope"].get("OPENVIKING_API_KEY") == "key-of-b" + + +@pytest.mark.asyncio +async def test_in_turn_teardown_keeps_the_callers_scope(tmp_path, monkeypatch): + monkeypatch.setenv("HERMES_HOME", str(tmp_path)) + seen: dict = {} + secret_scope.set_multiplex_active(True) + token = secret_scope.set_secret_scope({"MARKER": "turn-scope"}) + try: + await _runner({})._cleanup_agent_resources_off_loop(_recording_agent(seen), context="session hygiene") + finally: + secret_scope.reset_secret_scope(token) + secret_scope.set_multiplex_active(False) + assert seen["scope"] == {"MARKER": "turn-scope"}