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.
This commit is contained in:
@@ -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).
|
||||
|
||||
@@ -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()
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
@@ -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()
|
||||
|
||||
|
||||
Reference in New Issue
Block a user