diff --git a/agent/plugin_stream_hooks.py b/agent/plugin_stream_hooks.py index 0cb5d6a4d2..78ecd15448 100644 --- a/agent/plugin_stream_hooks.py +++ b/agent/plugin_stream_hooks.py @@ -8,6 +8,7 @@ registered are stopped lazily on the next lookup. from __future__ import annotations +import contextvars import logging import queue import threading @@ -26,7 +27,7 @@ _STOP = object() class _ConsumerDispatcher: hook_name: str callback: Callable[..., Any] - events: "queue.Queue[dict[str, Any] | object]" + events: "queue.Queue[tuple[contextvars.Context, dict[str, Any]] | object]" thread: threading.Thread | None = None @@ -62,10 +63,16 @@ def _worker(dispatcher: _ConsumerDispatcher) -> None: try: if item is _STOP: return - payload = dict(item) + context, payload = item + payload = dict(payload) payload.setdefault("telemetry_schema_version", OBSERVER_SCHEMA_VERSION) try: - dispatcher.callback(**payload) + # The worker outlives every turn and serves every profile; run the callback in the + # enqueuing turn's contextvars so it sees that turn's profile scope (home override, + # secrets), not an unbound context that fail-closed plugin bindings refuse (#118538). + # ``copy()``: one snapshot fans out to N consumer threads and a Context can only be + # entered by one thread at a time. + context.copy().run(dispatcher.callback, **payload) except Exception as exc: # Fires once per streaming delta: a mis-declared callback fails identically every # time, so it goes through the manager's warn-once reporter (#111922). @@ -126,7 +133,7 @@ def _dispatchers_for(hook_name: str) -> list[_ConsumerDispatcher]: def enqueue_plugin_stream_hook(hook_name: str, **payload: Any) -> bool: """Queue an observer hook for each consumer without running plugin code inline.""" queued = False - item = dict(payload) + item = (contextvars.copy_context(), dict(payload)) for dispatcher in _dispatchers_for(hook_name): if _put_drop_oldest(dispatcher.events, item): queued = True diff --git a/gateway/run.py b/gateway/run.py index 376873ad48..0c75256565 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -4896,8 +4896,9 @@ async def _shutdown_mcp_servers_nonblocking(timeout: float = 5.0, config: Any = the trailing wildcard pass — the only one that stops the shared loop — never ran. The worker runs in a FRESH context, not ``copy_context()``: the caller may sit inside a served - profile's scope, and ``launch_profile_scope_if_multiplexed`` documents "no HERMES_HOME override" - — inheriting one made the wildcard pass resolve the live home to that profile. + profile's scope, and the trailing wildcard pass must run under the launch profile's own scope + (``launch_profile_scope_if_multiplexed`` binds the launch home) — inheriting the caller's made it + resolve the live home to that profile. See #82874. """ diff --git a/hermes_cli/plugins_dispatch.py b/hermes_cli/plugins_dispatch.py index 3f54d05d2b..55b469fdb5 100644 --- a/hermes_cli/plugins_dispatch.py +++ b/hermes_cli/plugins_dispatch.py @@ -131,6 +131,9 @@ class _QueuedPluginEvent: subscriptions: tuple[_EventSubscription, ...] depth: int generation: int + # The emitter's contextvars: the single worker thread serves every profile, so each delivery + # runs under the profile scope the emit happened in (#118538). + context: contextvars.Context # Hook callback timeout (non-blocking abandon). Default cap per Python hook callback; overridden by @@ -396,7 +399,8 @@ class PluginDispatchMixin: callback = subscription.callback try: # Fresh deep copy per subscriber: no callback can mutate what the next sees. - resolve_plugin_command_result(callback(**copy.deepcopy(item.payload))) + resolve_plugin_command_result( + item.context.copy().run(callback, **copy.deepcopy(item.payload))) except (Exception, SystemExit) as exc: # A subscriber that fails identically on every emit is reported once (#111922). self._report_hook_failure(item.event, callback, item.payload, exc, surface="Event") @@ -431,7 +435,7 @@ class PluginDispatchMixin: return 0 item = _QueuedPluginEvent( event=event, payload=dict(payload), subscriptions=subscriptions, depth=depth + 1, - generation=generation) + generation=generation, context=contextvars.copy_context()) try: self._event_queue.put_nowait(item) except queue.Full: diff --git a/tests/tui_gateway/test_launch_profile_scope_plugin_hooks.py b/tests/tui_gateway/test_launch_profile_scope_plugin_hooks.py new file mode 100644 index 0000000000..67f172a5a6 --- /dev/null +++ b/tests/tui_gateway/test_launch_profile_scope_plugin_hooks.py @@ -0,0 +1,140 @@ +"""Plugin hooks fired from a launch-profile turn see a bound profile scope under multiplexing. + +Regression for #118538: routed turns bind ``_profile_runtime_scope(home)`` (HERMES_HOME override + +secrets + terminal policy), but the launch profile's own turns went through +``launch_profile_runtime_scope`` / ``_profile_runtime_scope_tokens(None)`` which bound secrets and +terminal policy only. Under multiplexing an unset override is the fail-closed "unbound context" +signal (``serves_routed_profile``, per-home slots, third-party runtime bindings such as OMH's +``pre_tool_call`` gate), so every plugin hook a launch-profile turn fired looked unscoped and a +fail-closed plugin vetoed every tool call. The dispatcher (``plugins_dispatch``) and the tool +executor already propagate contextvars; the seam is the launch scope itself. + +The two hook kinds delivered off-turn by long-lived worker threads (stream observers in +``agent.plugin_stream_hooks``, the plugin event bus in ``plugins_dispatch``) had the sibling gap: +the worker ran callbacks in its own empty context, so even a routed turn's observer saw no scope. +""" + +from __future__ import annotations + +from pathlib import Path + +import pytest + +from agent.secret_scope import set_multiplex_active +from hermes_cli import plugins as plugins_mod +from hermes_constants import get_hermes_home, get_hermes_home_override +from tools.daemon_pool import DaemonThreadPoolExecutor +from tools.thread_context import propagate_context_to_thread +import tui_gateway.server as server +from tui_gateway import launch_profile_policy as lpp + + +@pytest.fixture +def two_homes(tmp_path, monkeypatch): + launch = tmp_path / "hermes_home" + routed = launch / "profiles" / "beta" + for home, tag in ((launch, "alpha"), (routed, "beta")): + home.mkdir(parents=True) + (home / "config.yaml").write_text( + f"plugins:\n enabled: [stub]\n entries:\n stub:\n settings:\n x: {tag}\n", + encoding="utf-8") + (home / ".env").write_text("", encoding="utf-8") + monkeypatch.setenv("HERMES_HOME", str(launch)) + monkeypatch.setattr(server, "_hermes_home", launch) + monkeypatch.setattr(server, "_served_profile_homes", set()) + monkeypatch.setattr(lpp, "_snapshot", None) + plugins_mod._reset_plugin_managers_for_tests() + seen: list[dict] = [] + + def stub_pre_tool_call(**_kw): + from hermes_cli.config import load_config_readonly + entry = ((load_config_readonly().get("plugins") or {}).get("entries") or {}).get("stub") or {} + seen.append({"home": get_hermes_home().name, "x": (entry.get("settings") or {}).get("x"), + "bound": get_hermes_home_override() is not None}) + return None + + # Plugin managers are keyed per home: each profile loads its own copy of the plugin. + from hermes_constants import reset_hermes_home_override, set_hermes_home_override + for home in (launch, routed): + token = set_hermes_home_override(str(home)) + try: + manager = plugins_mod.get_plugin_manager() + finally: + reset_hermes_home_override(token) + manager._discovered = True # never scan the real plugin tree + manager._hooks.setdefault("pre_tool_call", []).append(stub_pre_tool_call) + yield launch, routed, seen + set_multiplex_active(False) + plugins_mod._reset_plugin_managers_for_tests() + + +def _fire_from_tool_worker(): + """The real hop: a tool worker thread + the bounded hook dispatcher.""" + pool = DaemonThreadPoolExecutor(max_workers=1) + try: + return pool.submit(propagate_context_to_thread( + lambda: plugins_mod._dispatch_pre_tool_call_hooks( + "terminal", {"command": "true"}, tool_call_id="c1", turn_id="t1"))).result(10) + finally: + pool.shutdown(wait=False) + + +def _fire_off_turn_workers(launch_manager_hook_seen): + """The two worker-thread deliveries: stream observer queue and the plugin event bus.""" + from agent import plugin_stream_hooks + manager = plugins_mod.get_plugin_manager() + observed: list[dict] = [] + + def observer(**_kw): + observed.append({"home": get_hermes_home().name, "bound": get_hermes_home_override() is not None}) + + manager._hooks.setdefault("on_stream_end", []).append(observer) + manager._subscribe_event("stub", "stub:tick", observer) + try: + assert plugin_stream_hooks.enqueue_plugin_stream_hook("on_stream_end", session_id="s") + assert manager._dispatch_event("stub:tick", {}) == 1 + assert manager._wait_for_event_dispatch(timeout=5.0) + plugin_stream_hooks.shutdown_plugin_stream_hook_dispatcher(timeout=5.0) + finally: + manager._hooks["on_stream_end"].remove(observer) + manager._remove_plugin_subscriptions("stub") + return observed + + +def test_launch_profile_turn_hooks_see_bound_scope_a_b_a(two_homes): + launch, routed, seen = two_homes + set_multiplex_active(True) + from gateway.run import _profile_runtime_scope + + with lpp.launch_profile_runtime_scope(launch): + _fire_from_tool_worker() + off_turn_a = _fire_off_turn_workers(seen) + with _profile_runtime_scope(routed, {}): + _fire_from_tool_worker() + off_turn_b = _fire_off_turn_workers(seen) + scopes = server._profile_runtime_scope_tokens(None) # serve backend: launch-profile session + try: + _fire_from_tool_worker() + finally: + server._release_profile_runtime_scope_tokens(scopes) + + assert [(r["home"], r["x"]) for r in seen] == [ + ("hermes_home", "alpha"), ("beta", "beta"), ("hermes_home", "alpha")] + assert all(r["bound"] for r in seen), seen + assert off_turn_a == [{"home": "hermes_home", "bound": True}] * 2, off_turn_a + assert off_turn_b == [{"home": "beta", "bound": True}] * 2, off_turn_b + assert get_hermes_home_override() is None # every scope released + + +def test_single_profile_launch_scope_binds_no_override(two_homes): + """Control: a host that never multiplexes keeps the standalone shape (no override, ``os.environ`` + precedence), so single-profile installs are byte-identical.""" + launch, _routed, seen = two_homes + scopes = server._profile_runtime_scope_tokens(None) + try: + assert get_hermes_home_override() is None + _fire_from_tool_worker() + finally: + server._release_profile_runtime_scope_tokens(scopes) + assert seen == [{"home": "hermes_home", "x": "alpha", "bound": False}] + assert Path(get_hermes_home()) == launch diff --git a/tui_gateway/launch_profile_policy.py b/tui_gateway/launch_profile_policy.py index 53bc2650aa..fd44989859 100644 --- a/tui_gateway/launch_profile_policy.py +++ b/tui_gateway/launch_profile_policy.py @@ -150,14 +150,23 @@ def launch_secret_scope(launch_home: "str | Path") -> Dict[str, str]: @contextlib.contextmanager def launch_profile_runtime_scope(launch_home: "str | Path") -> Iterator[None]: - """Bind the launch profile's own runtime scope for one body: ``launch_secret_scope`` plus its - terminal policy over the frozen launch ``TERMINAL_*`` overlay. No HERMES_HOME override — the - launch home IS the process home. For hosts whose launch-profile bodies are not RPC sessions - (the standalone messaging gateway after a hosted room activated multiplexing, #112878).""" + """Bind the launch profile's own runtime scope for one body: HERMES_HOME override naming the + launch home, ``launch_secret_scope``, and its terminal policy over the frozen launch + ``TERMINAL_*`` overlay. For hosts whose launch-profile bodies are not RPC sessions (the + standalone messaging gateway after a hosted room activated multiplexing, #112878). + + The home override is bound even though the launch home IS the process home: under multiplexing + "override unset" is the fail-closed signal for an UNBOUND context (``serves_routed_profile``, + plugin runtime bindings such as OMH's ``pre_tool_call`` gate, per-home slots keyed on the + override), so a launch-profile turn without it was indistinguishable from no turn at all and + every plugin hook it fired saw an unscoped process (#118538). Routed turns already bind theirs + (``gateway/run.py::_profile_runtime_scope``); the launch profile is a tenant like any other.""" from agent.secret_scope import reset_secret_scope, set_secret_scope + from hermes_constants import reset_hermes_home_override, set_hermes_home_override from tools.terminal_scope import install_profile_terminal_scope, reset_terminal_scope home = Path(launch_home) + home_token = set_hermes_home_override(str(home)) secret_token = set_secret_scope(launch_secret_scope(home)) terminal_token = install_profile_terminal_scope(home, env_overlay=launch_terminal_env()) try: @@ -165,6 +174,7 @@ def launch_profile_runtime_scope(launch_home: "str | Path") -> Iterator[None]: finally: reset_terminal_scope(terminal_token) reset_secret_scope(secret_token) + reset_hermes_home_override(home_token) def launch_profile_scope_if_multiplexed(): diff --git a/tui_gateway/model_switch.py b/tui_gateway/model_switch.py index 8e62abccdb..849702344e 100644 --- a/tui_gateway/model_switch.py +++ b/tui_gateway/model_switch.py @@ -78,14 +78,17 @@ def _profile_runtime_scope_tokens(profile_home, *, hydrate_secrets: bool = True) overlay = None scopes.home = set_hermes_home_override(str(home)) else: - # No home override: the launch home IS get_hermes_home() (``_profile_home`` answers None for - # "already the launch profile"); only its secrets (+ terminal policy under multiplex) need binding. + # The launch home IS get_hermes_home() (``_profile_home`` answers None for "already the + # launch profile"); single-profile, only its secrets need binding. Once multiplexing is + # active the override is bound too: an unset override is the "unbound context" signal + # plugin runtime bindings and per-home slots fail closed on (#118538). from tui_gateway.launch_profile_policy import launch_secret_scope, launch_terminal_env home = Path(_hermes_home) secrets = launch_secret_scope(home) scopes.secret = set_secret_scope(secrets) if not is_multiplex_active(): return scopes + scopes.home = set_hermes_home_override(str(home)) overlay = launch_terminal_env() if scopes.secret is None: scopes.secret = set_secret_scope(secrets)