fix(gateway): gateway-stop teardown fires memory-provider lifecycle hooks under the owning profile
Under `gateway.multiplex_profiles: true` every gateway-hosted session that ended
inside the gateway process logged `Memory provider 'openviking' on_session_end
failed: get_secret('OPENVIKING_API_KEY') called with no profile secret scope
active` and skipped its end-of-session commit. The provider lifecycle hooks
(`flush_pending` -> `on_session_end` -> provider teardown -> `close`) read
credentials and home at call time; the stop-time finalize pass and the
idle-cache sweep run on the main loop outside any adapter handler, so
`_run_in_executor_with_context` copied an EMPTY scope into the worker.
Route `_cleanup_agent_resources_off_loop` and `_finalize_session_off_loop`
through `_run_release_in_profile_scope` (the seam cache eviction already uses):
a scoped caller keeps its scope; an unscoped caller passes the session key and
the owner's profile scope is entered from it. The stop path passes the key for
both active and idle-cached agents. Cron's post-run cleanup thread is the
sibling seam and is fixed by the cherry-picked #110634.
Live repro (fake OpenViking recording the X-API-Key header, alpha profile,
multiplex on): base - WARNING + traceback, no commit; head - no warning,
`POST /api/v1/sessions/<sid>/commit` arrives with alpha's key.
Fixes #110622.
This commit is contained in:
+32
-11
@@ -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)
|
||||
|
||||
@@ -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(
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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"}
|
||||
Reference in New Issue
Block a user