From 8b181940b464a9db99f2a6256683f2e2d672514a Mon Sep 17 00:00:00 2001 From: teknium1 <127238744+teknium1@users.noreply.github.com> Date: Mon, 14 Sep 2026 23:31:08 -0700 Subject: [PATCH] fix(multiplex): background threads and teardown paths carry the turn's profile scope The profile scope (HERMES_HOME override, secret scope, terminal policy) is a contextvar bundle bound per turn. A bare threading.Thread / Timer / gRPC callback starts with an empty context and resolves the LAUNCH profile: - agent/title_generator.py: the auto-title thread read auxiliary.title_generation (model, language, provider key) from the default profile's config and billed the default's key for a secondary's session. Spawn via agent.memory_provider.spawn_context_thread (copy_context). - tui_gateway/session_lifecycle.py: every teardown caller is a bare Timer (ws-orphan reap), the idle-reaper thread, atexit _shutdown_sessions, the session.close pool RPC, superseded_by_resume or compute_host flush - none carries a scope, yet on_session_end / commit_memory_session / agent.close -> shutdown_memory_provider read the provider's config + credentials at call time. Under multiplex they failed closed (tail never committed, #110622 class); on the Desktop backend a secondary's transcript went to the launch profile's memory tenant. _finalize_session and _teardown_session now bind _session_profile_runtime_scope(session) around those blocks, which covers every spawn site through the single chokepoint. - plugins/platforms/google_chat/adapter.py: Pub/Sub callbacks run on the gRPC SubscriberClient's threads and run_coroutine_threadsafe copies THAT empty context onto the loop task, so _dispatch_message and everything under it (attachment cache, per-user OAuth token store via _acquire_user_chat_api -> _load_per_user_chat_api, TTS keys, delivery ledger, bot-id cache) resolved the launch profile. connect() captures its scope; _on_pubsub_message and _submit_on_loop run under a per-callback copy of it. spawn_context_thread gains a kwargs passthrough for the title thread's callbacks. --- agent/memory_provider.py | 4 +- agent/title_generator.py | 14 +-- plugins/platforms/google_chat/adapter.py | 26 ++++- .../test_background_thread_profile_scope.py | 99 +++++++++++++++++++ tui_gateway/session_lifecycle.py | 31 +++--- 5 files changed, 150 insertions(+), 24 deletions(-) create mode 100644 tests/tui_gateway/test_background_thread_profile_scope.py diff --git a/agent/memory_provider.py b/agent/memory_provider.py index b5a9d2edc5..938fd0b003 100644 --- a/agent/memory_provider.py +++ b/agent/memory_provider.py @@ -27,10 +27,10 @@ def ctx_bound(fn: Callable[..., Any]) -> Callable[..., Any]: def spawn_context_thread(target: Callable[..., Any], *, name: str, daemon: bool = True, - args: tuple = ()) -> threading.Thread: + args: tuple = (), kwargs: Optional[Dict[str, Any]] = None) -> threading.Thread: """Unstarted thread running *target* under the spawner's contextvars (see :func:`ctx_bound`). Every memory-provider background job (prefetch, sync, writer loops) must go through this.""" - return threading.Thread(target=ctx_bound(target), args=args, name=name, daemon=daemon) + return threading.Thread(target=ctx_bound(target), args=args, kwargs=kwargs, name=name, daemon=daemon) # v1 = best-effort on_pre_compress() with the raw message list; v2 = opt-in fail-closed # checkpoint (normalized evidence handoff + strict-mode failure propagation). diff --git a/agent/title_generator.py b/agent/title_generator.py index 19e1b5584c..341f2171cd 100644 --- a/agent/title_generator.py +++ b/agent/title_generator.py @@ -8,7 +8,6 @@ and neither replaces a name the user typed.""" import json import logging import re -import threading from contextlib import suppress from typing import Any, Callable, Optional @@ -488,10 +487,13 @@ def maybe_auto_title( logger.debug("Auto-title skipped: auxiliary.title_generation.enabled=false") return apply_instant_title(session_db, session_id, user_message, title_callback) - threading.Thread( - target=auto_title_session, + # The thread must resolve auxiliary.title_generation (config, provider key, language) for the + # profile whose turn this is: a bare Thread starts with an empty context and lands on the launch + # profile under multiplex, titling X's session with the default profile's model and billing its key. + from agent.memory_provider import spawn_context_thread + spawn_context_thread( + auto_title_session, name="auto-title", args=(session_db, session_id, user_message), - kwargs=dict(failure_callback=failure_callback, main_runtime=main_runtime, title_callback=title_callback, runtime_validator=runtime_validator), - daemon=True, - name="auto-title", + kwargs=dict(failure_callback=failure_callback, main_runtime=main_runtime, title_callback=title_callback, + runtime_validator=runtime_validator), ).start() diff --git a/plugins/platforms/google_chat/adapter.py b/plugins/platforms/google_chat/adapter.py index 46e29dae82..30c3453116 100644 --- a/plugins/platforms/google_chat/adapter.py +++ b/plugins/platforms/google_chat/adapter.py @@ -11,6 +11,7 @@ from __future__ import annotations import asyncio import contextlib +import contextvars import importlib import json import logging @@ -384,6 +385,12 @@ class GoogleChatAdapter(BasePlatformAdapter): self._project_id = self._subscription_path = self._bot_user_id = None # bot id is users/{id} self._supervisor_task: Optional[asyncio.Task] = None self._loop: Optional[asyncio.AbstractEventLoop] = None + # The profile scope this adapter was connected under (multiplex: HERMES_HOME override + secret + # scope). Pub/Sub callbacks arrive on the gRPC SubscriberClient's own threads with an EMPTY + # context, and ``run_coroutine_threadsafe`` copies THAT context onto the loop task — so + # ``_dispatch_message`` and everything it reaches (attachment cache, per-user OAuth token + # store, delivery ledger, TTS keys) would resolve the launch profile. Captured in ``connect()``. + self._scope_ctx: Optional[contextvars.Context] = None # User-authed Chat clients for native ``media.upload`` (bot identity is rejected # there) keyed by sender email; ``_user_credentials``/``_user_chat_api`` = LEGACY fallback. self._user_chat_api = self._user_credentials = None @@ -496,9 +503,13 @@ class GoogleChatAdapter(BasePlatformAdapter): return try: from agent.async_utils import safe_schedule_threadsafe - future = safe_schedule_threadsafe( - coro, loop, logger=logger, log_message="[GoogleChat] Failed to schedule background callback", - log_level=logging.WARNING, + # run_coroutine_threadsafe copies the CALLING thread's context onto the loop task; from the + # gRPC callback thread that is empty. Run the scheduling inside the adapter's connect-time + # scope so the task (and every to_thread/create_task under it) carries the profile. + ctx = self._scope_ctx.copy() if self._scope_ctx is not None else contextvars.copy_context() + future = ctx.run( + safe_schedule_threadsafe, coro, loop, logger=logger, + log_message="[GoogleChat] Failed to schedule background callback", log_level=logging.WARNING, ) except RuntimeError: logger.warning("[GoogleChat] Loop closed between check and submit") @@ -608,6 +619,7 @@ class GoogleChatAdapter(BasePlatformAdapter): message="google-cloud-pubsub / google-api-python-client not installed") return False self._loop = asyncio.get_running_loop() + self._scope_ctx = contextvars.copy_context() # the profile scope connect() runs under (see __init__) try: project_id, subscription_path = self._validate_config() credentials = self._load_sa_credentials() @@ -782,7 +794,13 @@ class GoogleChatAdapter(BasePlatformAdapter): def _on_pubsub_message(self, message: Any) -> None: """Pub/Sub callback — parse envelope and dispatch to the asyncio loop. Runs in a SubscriberClient worker thread: never block, never raise (that - triggers nack + infinite redelivery). Event type comes from ``ce-type``.""" + triggers nack + infinite redelivery). Event type comes from ``ce-type``. The body runs under + the adapter's profile scope (a per-callback copy: a Context cannot be entered concurrently) so + ``_save_cached_bot_id`` and the loop hand-off resolve the served profile, not the launch one.""" + ctx = self._scope_ctx.copy() if self._scope_ctx is not None else contextvars.copy_context() + ctx.run(self._handle_pubsub_message, message) + + def _handle_pubsub_message(self, message: Any) -> None: if self._shutting_down: message.nack() return diff --git a/tests/tui_gateway/test_background_thread_profile_scope.py b/tests/tui_gateway/test_background_thread_profile_scope.py new file mode 100644 index 0000000000..9ba30a5a7d --- /dev/null +++ b/tests/tui_gateway/test_background_thread_profile_scope.py @@ -0,0 +1,99 @@ +"""Background threads started from a profile-scoped turn carry that profile's scope. + +Under multiplex the HERMES_HOME override and secret scope are contextvars bound per turn; a bare +``threading.Thread`` / ``Timer`` starts with an EMPTY context and resolves the LAUNCH profile's config +and credentials (or fails closed on secrets). Two production spawn paths are exercised here through +their real entry points: the auto-title thread and the ws-orphan reap Timer → session teardown. +""" + +import threading +from pathlib import Path +from unittest.mock import MagicMock + +import pytest + +from agent.secret_scope import ( + UnscopedSecretError, build_profile_secret_scope, get_secret, reset_secret_scope, set_multiplex_active, + set_secret_scope, +) +from hermes_constants import get_hermes_home, reset_hermes_home_override, set_hermes_home_override + + +@pytest.fixture +def served_home(tmp_path, monkeypatch): + a = tmp_path / ".hermes" + b = a / "profiles" / "b" + b.mkdir(parents=True) + (a / ".env").write_text("TITLE_KEY=a-key\n", encoding="utf-8") + (b / ".env").write_text("TITLE_KEY=b-key\n", encoding="utf-8") + monkeypatch.setenv("HERMES_HOME", str(a)) + monkeypatch.setenv("TITLE_KEY", "a-key") + set_multiplex_active(True) + try: + yield a, b + finally: + set_multiplex_active(False) + + +def _observe_scope(seen: dict, done: threading.Event) -> None: + seen["home"] = str(get_hermes_home()) + try: + seen["key"] = get_secret("TITLE_KEY") + except UnscopedSecretError as exc: + seen["key"] = exc + done.set() + + +def test_auto_title_thread_runs_in_the_turns_profile_scope(served_home, monkeypatch): + """``maybe_auto_title`` fires ``auto_title_session`` on a thread; it must see the turn's profile.""" + import agent.title_generator as tg + + a, b = served_home + seen, done = {}, threading.Event() + monkeypatch.setattr(tg, "auto_title_session", lambda *args, **kwargs: _observe_scope(seen, done)) + monkeypatch.setattr(tg, "apply_instant_title", lambda *args, **kwargs: None) + db = MagicMock() + db.get_session_title.return_value = None + home_token = set_hermes_home_override(str(b)) + secret_token = set_secret_scope(build_profile_secret_scope(b)) + try: + tg.maybe_auto_title(db, "sess-1", "hello there", [{"role": "user", "content": "hello there"}]) + finally: + reset_secret_scope(secret_token) + reset_hermes_home_override(home_token) + assert done.wait(timeout=10), "auto-title thread never ran" + assert seen == {"home": str(b), "key": "b-key"} + + +def test_ws_orphan_reap_tears_down_under_the_sessions_profile(served_home, monkeypatch): + """The reap Timer (empty context) → ``_teardown_popped_session``: memory commit + ``agent.close`` run + under ``session['profile_home']``, so the provider reads B's config/credentials, not the launch's.""" + import tui_gateway.server as server + + a, b = served_home + seen_commit, seen_close = {}, {} + done = threading.Event() + agent = MagicMock() + agent._session_messages = None + agent.session_id = "sess-b" + agent.commit_memory_session = lambda history: _observe_scope(seen_commit, threading.Event()) + agent.close = lambda: _observe_scope(seen_close, done) + sid = "ws-b" + session = { + "agent": agent, "history": [{"role": "user", "content": "x"}], "history_lock": threading.Lock(), + "session_key": "sess-b", "profile_home": str(b), "transport": server._detached_ws_transport, + "last_active": 0.0, + } + monkeypatch.setattr(server, "_WS_ORPHAN_REAP_GRACE_S", 0.01) + monkeypatch.setattr(server, "_ws_session_is_orphaned", lambda s: s is session) + with server._sessions_lock: + server._sessions[sid] = session + try: + server._schedule_ws_orphan_reap(sid) + assert done.wait(timeout=10), "reap timer never tore the session down" + finally: + with server._sessions_lock: + server._sessions.pop(sid, None) + assert seen_commit == {"home": str(b), "key": "b-key"} + assert seen_close == {"home": str(b), "key": "b-key"} + assert Path(get_hermes_home()) == a # the test thread's own context is untouched diff --git a/tui_gateway/session_lifecycle.py b/tui_gateway/session_lifecycle.py index 5c385df9e3..dda33e2dd8 100644 --- a/tui_gateway/session_lifecycle.py +++ b/tui_gateway/session_lifecycle.py @@ -232,17 +232,22 @@ def _finalize_session(session: dict | None, end_reason: str = "tui_close") -> No if hasattr(agent, "_persist_session") and (snapshot := getattr(agent, "_session_messages", None)): with contextlib.suppress(Exception): agent._persist_session(snapshot) - # interrupted=True so crash-recovery plugins can flush state (mirrors cli.py atexit). - if agent is not None: - with contextlib.suppress(Exception): - from hermes_cli.lifecycle import invoke_hook - invoke_hook( - "on_session_end", completed=False, interrupted=True, - session_id=getattr(agent, "session_id", None) or session.get("session_key", ""), - model=getattr(agent, "model", "unknown"), platform=getattr(agent, "platform", None) or "tui") - if agent is not None and history and hasattr(agent, "commit_memory_session"): - with contextlib.suppress(Exception): - agent.commit_memory_session(history) + # interrupted=True so crash-recovery plugins can flush state (mirrors cli.py atexit). The end-of-session + # hooks and the memory commit read the provider's config/credentials at call time; every caller of this + # chokepoint is an unscoped reaper/Timer/atexit/pool thread, so bind the SESSION's profile here — unscoped + # they fail closed under multiplex (tail never committed) or, on the Desktop backend serving a named + # profile, commit a secondary's transcript to the launch profile's memory tenant (same class as #110622). + with _session_profile_runtime_scope(session): + if agent is not None: + with contextlib.suppress(Exception): + from hermes_cli.lifecycle import invoke_hook + invoke_hook( + "on_session_end", completed=False, interrupted=True, + session_id=getattr(agent, "session_id", None) or session.get("session_key", ""), + model=getattr(agent, "model", "unknown"), platform=getattr(agent, "platform", None) or "tui") + if agent is not None and history and hasattr(agent, "commit_memory_session"): + with contextlib.suppress(Exception): + agent.commit_memory_session(history) session_key = session.get("session_key") session_id = getattr(agent, "session_id", None) or session_key @@ -311,7 +316,9 @@ def _teardown_session(session: dict | None, *, end_reason: str = "tui_close") -> from tools.approval import unregister_gateway_notify if key := session.get("session_key"): unregister_gateway_notify(key) - with contextlib.suppress(Exception): + # agent.close() → shutdown_memory_provider reads the provider's config/credentials at call time; same + # scope rule as _finalize_session (every caller here is an unscoped reaper/atexit/pool thread). + with contextlib.suppress(Exception), _session_profile_runtime_scope(session): if hasattr(agent := session.get("agent"), "close"): agent.close()