From 208bd0b65a2b746cd755e375d209b8fa6b7fa987 Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Sat, 12 Sep 2026 01:03:46 -0700 Subject: [PATCH] fix(multiplex): per-profile plugin caches and the Yuanbao active adapter Ramp Router efforts cache + warm/disk flags, xAI and OpenRouter image catalogs, Hindsight append-capability verdict, memory-provider skill registry, OpenViking atexit provider, Honcho loopback flow status, Langfuse client (os.environ-only credentials) and disk-cleanup's protected cron paths held one profile's credential- or home-derived value process-wide; YuanbaoAdapter._active_instance was last-connected-wins across profiles. Keyed by home key / credential fingerprint under an override, credentials read through the secret scope, warm threads run under copy_context(); unscoped module slots stay for the single-profile path and the existing monkeypatch tests. --- gateway/platforms/yuanbao.py | 26 +- plugins/disk-cleanup/disk_cleanup.py | 8 +- plugins/image_gen/openrouter/__init__.py | 17 +- plugins/image_gen/xai/__init__.py | 40 ++- plugins/memory/__init__.py | 17 +- plugins/memory/hindsight/__init__.py | 17 +- plugins/memory/honcho/oauth_flow.py | 50 +++- plugins/memory/openviking/__init__.py | 37 ++- plugins/model-providers/router/__init__.py | 78 +++-- plugins/observability/langfuse/__init__.py | 76 +++-- .../test_multiplex_plugin_cache_scope.py | 272 ++++++++++++++++++ 11 files changed, 544 insertions(+), 94 deletions(-) create mode 100644 tests/plugins/test_multiplex_plugin_cache_scope.py diff --git a/gateway/platforms/yuanbao.py b/gateway/platforms/yuanbao.py index c8917f86ba..05a3c8e70e 100644 --- a/gateway/platforms/yuanbao.py +++ b/gateway/platforms/yuanbao.py @@ -2607,14 +2607,31 @@ class YuanbaoAdapter(BasePlatformAdapter): MEDIA_MAX_SIZE_MB: int = 50 DM_MAX_CHARS = 10000 _active_instance: ClassVar[Optional["YuanbaoAdapter"]] = None + # Per Hermes home: a multiplexed gateway runs one Yuanbao adapter per profile, and the tools / + # send_message read "the" adapter from inside a profile-scoped turn, so last-wins would route + # profile B's sends through profile A's bot. Registration and lookup both key on the ambient + # override (connect/reconnect tasks inherit the profile's Context); the slot above serves the + # unscoped path. + _active_instances: ClassVar[Dict[str, "YuanbaoAdapter"]] = {} @classmethod def get_active(cls) -> Optional["YuanbaoAdapter"]: - return cls._active_instance + from hermes_constants import get_hermes_home_override, hermes_home_key + + if get_hermes_home_override() is None: + return cls._active_instance + return cls._active_instances.get(hermes_home_key()) @classmethod def set_active(cls, adapter: Optional["YuanbaoAdapter"]) -> None: - cls._active_instance = adapter + from hermes_constants import get_hermes_home_override, hermes_home_key + + if get_hermes_home_override() is None: + cls._active_instance = adapter + elif adapter is None: + cls._active_instances.pop(hermes_home_key(), None) + else: + cls._active_instances[hermes_home_key()] = adapter def __init__(self, config: PlatformConfig, **kwargs: Any) -> None: super().__init__(config, Platform.YUANBAO) @@ -2697,7 +2714,10 @@ class YuanbaoAdapter(BasePlatformAdapter): async def disconnect(self) -> None: """Cancel background tasks and close the WebSocket connection.""" if YuanbaoAdapter._active_instance is self: - YuanbaoAdapter.set_active(None) + YuanbaoAdapter._active_instance = None + for home_key, active in list(YuanbaoAdapter._active_instances.items()): + if active is self: + del YuanbaoAdapter._active_instances[home_key] self._running = False self._mark_disconnected() self._release_platform_lock() diff --git a/plugins/disk-cleanup/disk_cleanup.py b/plugins/disk-cleanup/disk_cleanup.py index f45e27f982..a42479c690 100755 --- a/plugins/disk-cleanup/disk_cleanup.py +++ b/plugins/disk-cleanup/disk_cleanup.py @@ -102,20 +102,20 @@ _NEVER_TRACK_TOP_LEVEL = frozenset({ "patches", "projects", "skins", "themes", "contributors", "profiles", "backups", "optional-skills"}) -@functools.lru_cache(maxsize=1) # built lazily so HERMES_HOME resolves once -def _protected_cron_paths() -> frozenset: +@functools.lru_cache(maxsize=8) # keyed by home: a multiplexed process serves several profiles +def _protected_cron_paths(home: Path) -> frozenset: """Defense-in-depth for quick(): EXACT cron control-plane paths (``cron/``, ``output/`` root, ``jobs.json``, ``.tick.lock``) never deleted regardless of stored category (stale tracked.json). Never widen to everything under ``cron/output/``: run artifacts there are disposable; only wholesale deletion of ``output/`` is fatal.""" - return frozenset(str(x) for parent in ("cron", "cronjobs") for base in (get_hermes_home() / parent,) + return frozenset(str(x) for parent in ("cron", "cronjobs") for base in (home / parent,) for x in (base, base / "output", base / "jobs.json", base / ".tick.lock")) # Paths under $HERMES_HOME that must NEVER be deleted by quick(), regardless of what the stored category # says. This is a defense-in-depth guard against stale tracked.json entries from before #34840. def _is_protected_cron_path(p: Path) -> bool: - return str(p.resolve()) in _protected_cron_paths() + return str(p.resolve()) in _protected_cron_paths(get_hermes_home()) def fmt_size(n: float) -> str: diff --git a/plugins/image_gen/openrouter/__init__.py b/plugins/image_gen/openrouter/__init__.py index cc50385522..3c131bf9b5 100644 --- a/plugins/image_gen/openrouter/__init__.py +++ b/plugins/image_gen/openrouter/__init__.py @@ -62,9 +62,11 @@ _load_image_gen_config = load_image_gen_config _IMAGE_API_ENV_PREFIX = "OPENROUTER_IMAGE_API_" # Separate connect budget: no TLS in 20s means the endpoint is down — don't wait out the read budget. _IMAGE_API_CONNECT_TIMEOUT = 20.0 -# ``/images/models`` probes keyed by base URL: ``(fetched_at, ids)``; empty sets cached too. +# ``/images/models`` probes keyed by (base URL, key fingerprint): ``(fetched_at, ids)``; empty sets +# are cached too, so the key must include the credential or one profile's 401 would pin a sibling +# profile (same base URL, different key) to chat-completions for the whole TTL. _CATALOG_TTL_SECONDS = 900.0 -_CATALOG_CACHE: Dict[str, Tuple[float, frozenset]] = {} +_CATALOG_CACHE: Dict[Tuple[str, Optional[str]], Tuple[float, frozenset]] = {} _GEMINI_RATIOS = ( "1:1", "1:4", "1:8", "2:3", "3:2", "3:4", "4:1", "4:3", "4:5", "5:4", "8:1", "9:16", "16:9", "21:9", @@ -271,9 +273,12 @@ def _fetch_catalog( def _fetch_image_api_catalog(base_url: str, api_key: str) -> frozenset: - """Model ids from ``GET {base_url}/images/models``, cached per base URL. Any failure caches an - empty set (→ chat-completions): guessing "images" would 404 a working chat setup.""" - cached = _CATALOG_CACHE.get(base_url) + """Model ids from ``GET {base_url}/images/models``, cached per (base URL, key). Any failure caches + an empty set (→ chat-completions): guessing "images" would 404 a working chat setup.""" + from agent.credential_persistence import fingerprint_secret_value + + cache_key = (base_url, fingerprint_secret_value(api_key)) + cached = _CATALOG_CACHE.get(cache_key) if cached and (time.monotonic() - cached[0]) < _CATALOG_TTL_SECONDS: return cached[1] ids: set = set() @@ -283,7 +288,7 @@ def _fetch_image_api_catalog(base_url: str, api_key: str) -> frozenset: except Exception as exc: # noqa: BLE001 - probe must never break generation logger.debug("image API catalog probe failed for %s: %s", base_url, exc) resolved = frozenset(ids) - _CATALOG_CACHE[base_url] = (time.monotonic(), resolved) + _CATALOG_CACHE[cache_key] = (time.monotonic(), resolved) return resolved diff --git a/plugins/image_gen/xai/__init__.py b/plugins/image_gen/xai/__init__.py index 1bcd7b18a6..c6bf6e61b1 100644 --- a/plugins/image_gen/xai/__init__.py +++ b/plugins/image_gen/xai/__init__.py @@ -43,6 +43,10 @@ _EDIT_FALLBACK_MODEL = "grok-imagine-image-quality" # Live catalog cache ``(models, fetched_monotonic)``: ``/image-generation-models`` is the source of # truth (new models need no code change); ``_MODELS`` is the offline fallback + curated text. _LIVE_CACHE: Optional[Tuple[Dict[str, Dict[str, Any]], float]] = None +# Under a multiplexed profile override the catalog is keyed by (base_url, key fingerprint): the +# endpoint is credential-scoped, so one slot would hand profile A's models (or its cached auth +# failure) to profile B. The unscoped slot above stays for the single-profile path and its tests. +_LIVE_CACHE_BY_CREDENTIAL: Dict[Tuple[str, Optional[str]], Tuple[Dict[str, Dict[str, Any]], float]] = {} _LIVE_CACHE_TTL = 300.0 _LIVE_TIMEOUT = 10.0 @@ -63,9 +67,10 @@ def _base_url(creds: Dict[str, Any]) -> str: return str(creds.get("base_url") or "https://api.x.ai/v1").strip().rstrip("/") -def _fetch_live_models() -> Dict[str, Dict[str, Any]]: +def _fetch_live_models(creds: Optional[Dict[str, Any]] = None) -> Dict[str, Dict[str, Any]]: """``{model_id: {"input_modalities", "aliases"}}`` from the live endpoint; raises on failure.""" - creds = resolve_xai_http_credentials() + if creds is None: + creds = resolve_xai_http_credentials() api_key = str(creds.get("api_key") or "").strip() if not api_key: raise RuntimeError("no xAI credentials") @@ -88,15 +93,36 @@ def _fetch_live_models() -> Dict[str, Dict[str, Any]]: def _live_models() -> Dict[str, Dict[str, Any]]: """Cached live catalog (``{}`` when unreachable).""" global _LIVE_CACHE - if _LIVE_CACHE is not None and time.monotonic() - _LIVE_CACHE[1] < _LIVE_CACHE_TTL: + from hermes_constants import get_hermes_home_override + + if get_hermes_home_override() is None: + if _LIVE_CACHE is not None and time.monotonic() - _LIVE_CACHE[1] < _LIVE_CACHE_TTL: + return _LIVE_CACHE[0] + _LIVE_CACHE = (_fetch_live_models_or_empty(None), time.monotonic()) return _LIVE_CACHE[0] + + from agent.credential_persistence import fingerprint_secret_value + try: - live = _fetch_live_models() + creds = resolve_xai_http_credentials() + except Exception as exc: # noqa: BLE001 - unresolvable credentials → static fallback + logger.debug("xAI live image model catalog unavailable: %s", exc) + creds = {} + key = (_base_url(creds), fingerprint_secret_value(creds.get("api_key"))) + cached = _LIVE_CACHE_BY_CREDENTIAL.get(key) + if cached is not None and time.monotonic() - cached[1] < _LIVE_CACHE_TTL: + return cached[0] + live = _fetch_live_models_or_empty(creds) + _LIVE_CACHE_BY_CREDENTIAL[key] = (live, time.monotonic()) + return live + + +def _fetch_live_models_or_empty(creds: Optional[Dict[str, Any]]) -> Dict[str, Dict[str, Any]]: + try: + return _fetch_live_models() if creds is None else _fetch_live_models(creds) except Exception as exc: # noqa: BLE001 - offline/unauth → static fallback logger.debug("xAI live image model catalog unavailable: %s", exc) - live = {} - _LIVE_CACHE = (live, time.monotonic()) - return live + return {} def _catalog() -> Dict[str, Dict[str, Any]]: diff --git a/plugins/memory/__init__.py b/plugins/memory/__init__.py index 0f20bb0e50..7a9fae23f1 100644 --- a/plugins/memory/__init__.py +++ b/plugins/memory/__init__.py @@ -28,7 +28,15 @@ logger = logging.getLogger(__name__) _MEMORY_PLUGINS_DIR = Path(__file__).parent ENTRY_POINTS_GROUP = "hermes_agent.memory_providers" -_REGISTERED_MEMORY_PROVIDER_SKILLS: dict[str, Path] = {} +# Per Hermes home (plugin managers are per home too): pruning under one multiplexed profile must +# only retract that profile's provider skills, never a sibling profile's. +_REGISTERED_MEMORY_PROVIDER_SKILLS: dict[str, dict[str, Path]] = {} + + +def _registered_skills_for_active_home() -> dict[str, Path]: + from hermes_constants import hermes_home_key + + return _REGISTERED_MEMORY_PROVIDER_SKILLS.setdefault(hermes_home_key(), {}) # Synthetic parent package so user-installed providers don't collide with bundled ones. _USER_NAMESPACE = "_hermes_user_memory" @@ -330,7 +338,7 @@ class _ProviderCollector: registered_path = get_plugin_manager().find_plugin_skill(qualified_name) if registered_path is not None: - _REGISTERED_MEMORY_PROVIDER_SKILLS[qualified_name] = registered_path + _registered_skills_for_active_home()[qualified_name] = registered_path except Exception as exc: logger.debug("Memory provider '%s' failed to register skill: %s", self.name, exc) @@ -383,12 +391,13 @@ def _prune_inactive_memory_provider_skills(active_provider: Optional[str] = None from hermes_cli.plugins import get_plugin_manager manager = get_plugin_manager() - for qualified_name, registered_path in list(_REGISTERED_MEMORY_PROVIDER_SKILLS.items()): + registered = _registered_skills_for_active_home() + for qualified_name, registered_path in list(registered.items()): if qualified_name.partition(":")[0] == active_provider: continue if manager.find_plugin_skill(qualified_name) == registered_path: manager.remove_plugin_skill(qualified_name) - _REGISTERED_MEMORY_PROVIDER_SKILLS.pop(qualified_name, None) + registered.pop(qualified_name, None) def discover_plugin_cli_commands() -> List[dict]: diff --git a/plugins/memory/hindsight/__init__.py b/plugins/memory/hindsight/__init__.py index 381b43c657..25e663e761 100644 --- a/plugins/memory/hindsight/__init__.py +++ b/plugins/memory/hindsight/__init__.py @@ -91,9 +91,11 @@ def _maybe_upgrade_client() -> None: pass # packaging not available or other issue — proceed anyway -# update_mode='append' capability (Hindsight >= 0.5.0), cached per API URL per -# process so every provider on the same API shares one /version round trip. -_append_capability_cache: Dict[str, bool] = {} +# update_mode='append' capability (Hindsight >= 0.5.0), cached per (API URL, key fingerprint) +# per process so every provider on the same API+key shares one /version round trip. A failed probe +# caches False, so the key must include the credential or one profile's 401 would silently downgrade +# a sibling profile that shares the URL with a valid key. +_append_capability_cache: Dict[tuple[str, str | None], bool] = {} _append_capability_lock = threading.Lock() @@ -123,9 +125,12 @@ def _check_api_supports_update_mode_append(api_url: str, api_key: str | None = N """ if not api_url: return False + from agent.credential_persistence import fingerprint_secret_value + + cache_key = (api_url, fingerprint_secret_value(api_key)) with _append_capability_lock: - if api_url in _append_capability_cache: - return _append_capability_cache[api_url] + if cache_key in _append_capability_cache: + return _append_capability_cache[cache_key] version = _fetch_hindsight_api_version(api_url, api_key) try: # missing/invalid version -> unsupported from packaging.version import Version @@ -134,7 +139,7 @@ def _check_api_supports_update_mode_append(api_url: str, api_key: str | None = N supported = False with _append_capability_lock: # A concurrent probe may have filled the cache meanwhile; its answer wins. - supported = _append_capability_cache.setdefault(api_url, supported) + supported = _append_capability_cache.setdefault(cache_key, supported) if supported: logger.debug("Hindsight API %s version %s supports update_mode='append'", api_url, version) else: diff --git a/plugins/memory/honcho/oauth_flow.py b/plugins/memory/honcho/oauth_flow.py index f29ec7240a..bb89c11dfb 100644 --- a/plugins/memory/honcho/oauth_flow.py +++ b/plugins/memory/honcho/oauth_flow.py @@ -364,6 +364,25 @@ class FlowStatus: _status = FlowStatus() _status_lock = threading.Lock() _flow_thread: threading.Thread | None = None +# Status + thread per (config_path, host): the flow writes ONE host block of ONE honcho.json, so two +# profiles connecting in the same process must not share (or refuse each other on) one status slot. +# The module slots above serve the unscoped single-profile path (and its tests). +_flows_by_target: dict[tuple[str, str], tuple[FlowStatus, threading.Thread | None]] = {} + + +def _flow_target() -> tuple[str, str] | None: + """(config_path, host) of the active profile override, or None when unscoped.""" + from hermes_constants import get_hermes_home_override + + if get_hermes_home_override() is None: + return None + return str(resolve_config_path()), resolve_active_host() + + +def _flow_state(target: tuple[str, str] | None) -> tuple[FlowStatus, threading.Thread | None]: + if target is None: + return _status, _flow_thread + return _flows_by_target.setdefault(target, (FlowStatus(), None)) def _detect_connection() -> tuple[bool, str | None]: """Report whether a credential is already stored: 'oauth', 'apikey', or none.""" @@ -376,14 +395,15 @@ def _detect_connection() -> tuple[bool, str | None]: return auth is not None, auth def get_flow_status() -> dict[str, object]: + status, _thread = _flow_state(_flow_target()) with _status_lock: - state, detail = _status.state, _status.detail + state, detail = status.state, status.detail connected, auth = _detect_connection() return {"state": state, "detail": detail, "connected": connected, "auth": auth} -def _set_status(state: str, detail: str = "") -> None: +def _set_status(status: FlowStatus, state: str, detail: str = "") -> None: with _status_lock: - _status.state, _status.detail = state, detail + status.state, status.detail = state, detail def start_loopback_flow_background( *, config_path: Path | None = None, host: str | None = None, source: str = "hermes-desktop", @@ -393,21 +413,27 @@ def start_loopback_flow_background( Idempotent while pending, so a double-click can't open two tabs / bind :8765 twice.""" global _flow_thread # Resolve under the caller's profile scope NOW — a context-local HERMES_HOME override can't reach the worker. - config_path = config_path or resolve_config_path() - host = host or resolve_active_host() + target = _flow_target() + config_path = config_path or (Path(target[0]) if target else resolve_config_path()) + host = host or (target[1] if target else resolve_active_host()) + status, thread = _flow_state(target) with _status_lock: - if _status.state == "pending" and _flow_thread and _flow_thread.is_alive(): - return {"state": _status.state, "detail": _status.detail} - _status.state, _status.detail = "pending", "waiting for browser consent" + if status.state == "pending" and thread and thread.is_alive(): + return {"state": status.state, "detail": status.detail} + status.state, status.detail = "pending", "waiting for browser consent" def _run() -> None: try: authorize_via_loopback(config_path=config_path, host=host, source=source, timeout=timeout) - _set_status("connected", "Honcho connected") + _set_status(status, "connected", "Honcho connected") except Exception as exc: logger.warning("Honcho OAuth loopback flow failed: %s", exc) - _set_status("error", str(exc)) + _set_status(status, "error", str(exc)) - _flow_thread = threading.Thread(target=_run, name="honcho-oauth-loopback", daemon=True) - _flow_thread.start() + thread = threading.Thread(target=_run, name="honcho-oauth-loopback", daemon=True) + if target is None: + _flow_thread = thread + else: + _flows_by_target[target] = (status, thread) + thread.start() return get_flow_status() diff --git a/plugins/memory/openviking/__init__.py b/plugins/memory/openviking/__init__.py index f5a9ee7f22..f29341be05 100644 --- a/plugins/memory/openviking/__init__.py +++ b/plugins/memory/openviking/__init__.py @@ -198,23 +198,23 @@ def _preview(value: Any, limit: int = 160) -> str: # atexit safety net: commit pending sessions even if shutdown_memory_provider # never runs (gateway crash, exception in the session expiry watcher, ...). -_last_active_provider: Optional["OpenVikingMemoryProvider"] = None +# One entry per Hermes home: a multiplexed gateway initializes a provider per profile and every +# one of them holds pending sessions worth committing, not just the last to initialize. +_active_providers_by_home: Dict[str, "OpenVikingMemoryProvider"] = {} def _atexit_commit_sessions(): - global _last_active_provider - provider = _last_active_provider - if provider is None: - return - _last_active_provider = None - try: - with suppress(Exception): # best-effort at shutdown time - provider.on_session_end([]) - finally: - # ``finally`` (as on main): the run lock is released even when on_session_end - # dies of a BaseException (KeyboardInterrupt during atexit). - with suppress(Exception): - provider._release_run_lock() + providers = list(_active_providers_by_home.values()) + _active_providers_by_home.clear() + for provider in providers: + try: + with suppress(Exception): # best-effort at shutdown time + provider.on_session_end([]) + finally: + # ``finally`` (as on main): the run lock is released even when on_session_end + # dies of a BaseException (KeyboardInterrupt during atexit). + with suppress(Exception): + provider._release_run_lock() atexit.register(_atexit_commit_sessions) @@ -1436,8 +1436,7 @@ class OpenVikingMemoryProvider(MemoryProvider): self._conn_snapshot = self._settings_tuple() self._recover_pending_sessions() - global _last_active_provider # atexit safety net - _last_active_provider = self + _active_providers_by_home[self._hermes_home] = self # atexit safety net def _ensure_client(self) -> Optional["_VikingClient"]: """Active client, rebuilt if the resolved config changed. @@ -2488,9 +2487,9 @@ class OpenVikingMemoryProvider(MemoryProvider): for t in workers: if t.is_alive(): t.join(timeout=5.0) - global _last_active_provider # clear so atexit doesn't double-commit - if _last_active_provider is self: - _last_active_provider = None + # Clear so atexit doesn't double-commit. + if _active_providers_by_home.get(self._hermes_home) is self: + del _active_providers_by_home[self._hermes_home] self._release_run_lock() @staticmethod diff --git a/plugins/model-providers/router/__init__.py b/plugins/model-providers/router/__init__.py index 1c485e8d02..7dd072ef65 100644 --- a/plugins/model-providers/router/__init__.py +++ b/plugins/model-providers/router/__init__.py @@ -8,9 +8,11 @@ from ``GET /v1/models`` is cached (memory + disk mirror, background warmer; neve request hot path) and fed to the codex transport's clamp via ``supported_reasoning_efforts``. """ +import contextvars import json import logging import os +import sys import threading import time from pathlib import Path @@ -32,13 +34,49 @@ _efforts_lock = threading.Lock() _warm_started = False _disk_checked = False -# A stale verdict beats no verdict: a past-TTL mirror is still served while a -# background refresh runs. -_DISK_TTL_SECONDS = 24 * 60 * 60 + +class _CacheState: + """Efforts cache + once-only flags for one Hermes home (same names as the module slots).""" + + __slots__ = ("_efforts_cache", "_warm_started", "_disk_checked") + + def __init__(self) -> None: + self._efforts_cache: Optional[dict[str, list[str]]] = None + self._warm_started = False + self._disk_checked = False + + +# The catalog is account-scoped and the disk mirror lives under each profile's home, so under a +# multiplexed profile override the memory cache and its once-only flags are per home too; +# otherwise profile B's clamp would be built from profile A's key (and A's warm/disk flags). +_state_by_home: dict[str, _CacheState] = {} + + +def _state() -> Any: + """Holder of ``_efforts_cache``/``_warm_started``/``_disk_checked``: this module when unscoped + (tests monkeypatch those slots), else the active home's ``_CacheState``.""" + from hermes_constants import get_hermes_home_override, hermes_home_key + + if get_hermes_home_override() is None: + return sys.modules[__name__] + with _efforts_lock: + return _state_by_home.setdefault(hermes_home_key(), _CacheState()) def _base_url() -> str: - return os.getenv("RAMP_ROUTER_BASE_URL", "").strip().rstrip("/") or ROUTER_DEFAULT_BASE_URL + """Router base URL: profile ``.env`` first (scope-aware), plain os.environ as the fallback.""" + try: + from hermes_cli.config import get_env_value_prefer_dotenv as prefer_dotenv + except Exception: + prefer_dotenv = None + for resolve in filter(None, (prefer_dotenv, os.environ.get)): + try: + value = str(resolve("RAMP_ROUTER_BASE_URL") or "").strip().rstrip("/") + except Exception: + value = "" + if value: + return value + return ROUTER_DEFAULT_BASE_URL def _resolve_api_key() -> str: @@ -141,11 +179,11 @@ def _load_disk() -> tuple[Optional[dict[str, list[str]]], float]: def _seed_efforts(items: Any) -> Optional[dict[str, list[str]]]: """Seed memory + disk caches from a ``/v1/models`` payload.""" - global _efforts_cache parsed = _parse_efforts(items) if parsed is not None: + state = _state() with _efforts_lock: - _efforts_cache = parsed + state._efforts_cache = parsed _save_disk(parsed) return parsed @@ -173,36 +211,38 @@ def _fetch_catalog_items(*, api_key: str = "", base_url: str = "", timeout: floa def _efforts_cache_only() -> Optional[dict[str, list[str]]]: - """Memory, else the disk mirror (checked once per process). Never HTTP (hot-path safe).""" - global _efforts_cache, _disk_checked + """Memory, else the disk mirror (checked once per home). Never HTTP (hot-path safe).""" + state = _state() with _efforts_lock: - cached = _efforts_cache - if cached is not None or _disk_checked: + cached = state._efforts_cache + if cached is not None or state._disk_checked: return cached - _disk_checked = True + state._disk_checked = True parsed, age = _load_disk() if parsed is None: return None with _efforts_lock: - _efforts_cache = cached = _efforts_cache if _efforts_cache is not None else parsed + if state._efforts_cache is None: + state._efforts_cache = parsed + cached = state._efforts_cache if age >= _DISK_TTL_SECONDS: _warm_efforts_async() return cached def _warm_efforts_async() -> None: - """Refresh the efforts cache in the background, at most once per process. + """Refresh the efforts cache in the background, at most once per home. Skipped under pytest (a mid-suite fetch makes cache state timing-dependent) and without a key (it would 401; the first authenticated fetch_models() seeds). """ - global _warm_started if os.environ.get("PYTEST_CURRENT_TEST"): return + state = _state() with _efforts_lock: - if _warm_started: + if state._warm_started: return - _warm_started = True + state._warm_started = True if not _resolve_api_key(): return @@ -211,7 +251,11 @@ def _warm_efforts_async() -> None: if items is not None: _seed_efforts(items) try: - threading.Thread(target=_refresh, name="router-caps-warm", daemon=True).start() + # copy_context: the home override / secret scope are ContextVars, so a bare thread would + # fetch with the launch profile's key and mirror into its cache dir. + threading.Thread( + target=contextvars.copy_context().run, args=(_refresh,), name="router-caps-warm", daemon=True, + ).start() except Exception as exc: logger.debug("router: caps warmer failed to start: %s", exc) diff --git a/plugins/observability/langfuse/__init__.py b/plugins/observability/langfuse/__init__.py index a6edc25227..cd1c22e294 100644 --- a/plugins/observability/langfuse/__init__.py +++ b/plugins/observability/langfuse/__init__.py @@ -52,6 +52,10 @@ _TRACE_STATE: Dict[str, TraceState] = {} # Bounds the leak, not concurrency. _MAX_TRACE_STATE = 256 _LANGFUSE_CLIENT = None +# Under a multiplexed profile override, one settled client (or _INIT_FAILED) per Hermes home: the +# keys live in each profile's .env, so a single slot would trace profile B into profile A's project +# (or pin B to A's failed init). The slot above stays for the unscoped single-profile path. +_LANGFUSE_CLIENT_BY_HOME: Dict[str, Any] = {} # Separate from _STATE_LOCK (hot path) so the two never nest; serializes the # first client build so racing callers can't each construct a client. _LANGFUSE_CLIENT_LOCK = threading.Lock() @@ -83,6 +87,19 @@ def _env(name: str, default: str = "") -> str: return os.environ.get(name, default).strip() +def _secret(name: str) -> str: + """Credential read honoring the active profile's secret scope; plain os.environ when unscoped.""" + try: + from agent.secret_scope import UnscopedSecretError, get_secret + try: + return (get_secret(name) or "").strip() + except UnscopedSecretError: + pass + except Exception: + pass + return _env(name) + + def _debug(message: str) -> None: if _env("HERMES_LANGFUSE_DEBUG").lower() in {"1", "true", "yes", "on"}: logger.info("Langfuse tracing: %s", message) @@ -181,24 +198,48 @@ def _validate_langfuse_key(env_name: str, value: str) -> Optional[str]: return f"{env_name}={preview} (expected {expected!r} prefix)" +def _settled_client() -> Any: + """The active profile's settled client slot value (client, ``_INIT_FAILED`` or ``None`` = never + built). Never initializes.""" + from hermes_constants import get_hermes_home_override, hermes_home_key + + if get_hermes_home_override() is None: + return _LANGFUSE_CLIENT + return _LANGFUSE_CLIENT_BY_HOME.get(hermes_home_key()) + + +def _settle_client() -> Any: + """Build once and store for the active profile. Caller holds ``_LANGFUSE_CLIENT_LOCK``.""" + global _LANGFUSE_CLIENT + from hermes_constants import get_hermes_home_override, hermes_home_key + + client = _build_client() + settled = _INIT_FAILED if client is None else client + if get_hermes_home_override() is None: + _LANGFUSE_CLIENT = settled + else: + _LANGFUSE_CLIENT_BY_HOME[hermes_home_key()] = settled + if client is not None: + # atexit is LIFO: registering AFTER the SDK's constructor means our + # finalizer runs first, so root spans ended there still get flushed + # by the SDK (short-lived processes: kanban workers, chat -q, cron). + atexit.register(_finalize_all_traces) + return settled + + def _get_langfuse() -> Optional[Langfuse]: """Cached Langfuse client, or ``None`` if the SDK/credentials are unavailable. The first build is serialized so racing callers can't each construct a client and leak the loser's HTTP connection + flush thread.""" - global _LANGFUSE_CLIENT # Fast path — already settled (success or _INIT_FAILED) needs no lock; # re-check under it since a racing thread may have finished init. - if _LANGFUSE_CLIENT is None: + settled = _settled_client() + if settled is None: with _LANGFUSE_CLIENT_LOCK: - if _LANGFUSE_CLIENT is None: - client = _build_client() - _LANGFUSE_CLIENT = _INIT_FAILED if client is None else client - if client is not None: - # atexit is LIFO: registering AFTER the SDK's constructor means our - # finalizer runs first, so root spans ended there still get flushed - # by the SDK (short-lived processes: kanban workers, chat -q, cron). - atexit.register(_finalize_all_traces) - return None if _LANGFUSE_CLIENT is _INIT_FAILED else _LANGFUSE_CLIENT + settled = _settled_client() + if settled is None: + settled = _settle_client() + return None if settled is _INIT_FAILED else settled def _build_client() -> Optional[Langfuse]: @@ -211,7 +252,7 @@ def _build_client() -> Optional[Langfuse]: ) return None - public_key, secret_key = (_env(f"HERMES_LANGFUSE_{n}") or _env(f"LANGFUSE_{n}") for n in ("PUBLIC_KEY", "SECRET_KEY")) + public_key, secret_key = (_secret(f"HERMES_LANGFUSE_{n}") or _secret(f"LANGFUSE_{n}") for n in ("PUBLIC_KEY", "SECRET_KEY")) if not (public_key and secret_key): return None @@ -234,10 +275,10 @@ def _build_client() -> Optional[Langfuse]: kwargs: Dict[str, Any] = {"public_key": public_key, "secret_key": secret_key} for key, name, default in (("base_url", "BASE_URL", "https://cloud.langfuse.com"), ("environment", "ENV", ""), ("release", "RELEASE", "")): - value = _env(f"HERMES_LANGFUSE_{name}") or _env(f"LANGFUSE_{name}") or default + value = _secret(f"HERMES_LANGFUSE_{name}") or _secret(f"LANGFUSE_{name}") or default if value: kwargs[key] = value - sample_rate = _env("HERMES_LANGFUSE_SAMPLE_RATE") + sample_rate = _secret("HERMES_LANGFUSE_SAMPLE_RATE") if sample_rate: try: kwargs["sample_rate"] = float(sample_rate) @@ -596,7 +637,10 @@ def _finalize_all_traces() -> None: _end_children(state, include_subagents=True) _end_root(state, f"atexit finalize for {key}") if states: - _flush(_get_langfuse()) + # atexit runs unscoped; flush every profile's client, not just the launch profile's. + for client in (_get_langfuse(), *_LANGFUSE_CLIENT_BY_HOME.values()): + if client is not _INIT_FAILED: + _flush(client) def _flush(client: Any) -> None: @@ -909,7 +953,7 @@ def on_session_finalize(*, session_id: str = "", reason: str = "", **_: Any) -> tool-only or empty final response never reaches ``_finish_trace``; its root would dangle until eviction and queued events could be lost on exit.""" # Never lazily initialize a client here — if init never happened there are no traces. - client = _LANGFUSE_CLIENT + client = _settled_client() if client is None or client is _INIT_FAILED or not hasattr(client, "flush"): return diff --git a/tests/plugins/test_multiplex_plugin_cache_scope.py b/tests/plugins/test_multiplex_plugin_cache_scope.py new file mode 100644 index 0000000000..9b70532def --- /dev/null +++ b/tests/plugins/test_multiplex_plugin_cache_scope.py @@ -0,0 +1,272 @@ +"""Plugin module-level caches must not hand profile A's state to profile B under a multiplexed +HERMES_HOME override (``hermes_constants.set_hermes_home_override``). + +One invariant per mechanism: home-keyed slot with the unscoped module slot intact (router; yuanbao's +ClassVar twin), credential-fingerprinted catalog keys (openrouter), per-home registries (memory +provider skills), collect-all atexit (openviking), lru_cache keyed by the home (disk-cleanup). +Only HTTP transports are faked; the caches themselves are exercised for real. +""" + +from __future__ import annotations + +import contextlib +import importlib.util +import json +import sys +import threading +import time +from pathlib import Path + +import pytest + +from agent.secret_scope import build_profile_secret_scope, reset_secret_scope, set_secret_scope +from hermes_constants import reset_hermes_home_override, set_hermes_home_override + +REPO = Path(__file__).resolve().parents[2] + + +@pytest.fixture +def homes(tmp_path, monkeypatch): + """Profile A (launch home, ``HERMES_HOME``) and profile B with different config/.env values.""" + root = tmp_path / ".hermes" + a, b = root, root / "profiles" / "B" + for home, tag in ((a, "A"), (b, "B")): + home.mkdir(parents=True) + (home / "config.yaml").write_text(f"memory:\n provider: prov{tag}\n", encoding="utf-8") + (home / ".env").write_text( + f"RAMP_ROUTER_API_KEY=router-key-{tag}\nRAMP_ROUTER_BASE_URL=https://{tag.lower()}.router.test/v1\n", + encoding="utf-8") + monkeypatch.setenv("HERMES_HOME", str(a)) + for var in ("RAMP_ROUTER_API_KEY", "RAMP_ROUTER_BASE_URL", "PYTEST_CURRENT_TEST"): + monkeypatch.delenv(var, raising=False) + return a, b + + +@contextlib.contextmanager +def scoped(home: Path): + t_home = set_hermes_home_override(str(home)) + t_secret = set_secret_scope(build_profile_secret_scope(home)) + try: + yield + finally: + reset_secret_scope(t_secret) + reset_hermes_home_override(t_home) + + +class _Resp: + def __init__(self, payload, status=200): + self._payload, self.status_code = payload, status + + def read(self): + return json.dumps(self._payload).encode() + + def json(self): + return self._payload + + def raise_for_status(self): + if self.status_code >= 400: + raise RuntimeError(f"HTTP {self.status_code}") + + def __enter__(self): + return self + + def __exit__(self, *_exc): + return False + + +def _router(): + from providers import get_provider_profile + + profile = get_provider_profile("router") + return profile, sys.modules[type(profile).__module__] + + +def test_router_efforts_cache_and_base_url_follow_the_active_profile(homes, monkeypatch): + """Efforts map + once-only flags are per home under an override (and the warm thread inherits the + scope), while the unscoped path keeps using the module slots; the base URL comes from the + profile's .env.""" + import hermes_cli.urllib_security as urllib_security + + a, b = homes + profile, mod = _router() + fetched: list[str] = [] + + def fake_open(req, *, timeout, **_kw): + tag = (req.get_header("Authorization") or "").rsplit("-", 1)[-1] + fetched.append(req.full_url) + return _Resp({"data": [{"id": f"model-{tag}", "router": {"capabilities": {"reasoning": { + "supported": True, "efforts": [{"value": "low" if tag == "A" else "high"}]}}}}]}) + + monkeypatch.setattr(urllib_security, "open_credentialed_url", fake_open) + monkeypatch.setattr(mod, "_efforts_cache", None) + monkeypatch.setattr(mod, "_disk_checked", False) + monkeypatch.setattr(mod, "_warm_started", False) + + with scoped(a): + assert mod._base_url() == "https://a.router.test/v1" + profile.fetch_models() + assert profile.supported_reasoning_efforts("model-A") == ("low",) + with scoped(b): + assert mod._base_url() == "https://b.router.test/v1" + # B never fetched: A's verdicts must not be visible, and B's disk mirror (absent) is what it reads. + assert profile.supported_reasoning_efforts("model-A") is None + profile.fetch_models() + assert profile.supported_reasoning_efforts("model-B") == ("high",) + with scoped(a): + assert profile.supported_reasoning_efforts("model-A") == ("low",) + # Unscoped: the module slot is untouched by the scoped fetches. + assert mod._efforts_cache is None + + # Warm thread launched from B's turn fetches with B's key/base URL. pytest re-sets + # PYTEST_CURRENT_TEST per phase; the warmer's pytest guard reads it at call time. + fetched.clear() + monkeypatch.delenv("PYTEST_CURRENT_TEST", raising=False) + with scoped(b): + mod._warm_efforts_async() + deadline = time.monotonic() + 5 + while time.monotonic() < deadline and not fetched: + time.sleep(0.02) + assert fetched and all(url.startswith("https://b.router.test/") for url in fetched) + + +def test_credentialed_catalog_probe_failure_is_not_cached_across_keys(monkeypatch): + """A 401 under one key must not pin a sibling profile (same base URL, valid key) to the empty + catalog for the TTL.""" + import requests + + import plugins.image_gen.openrouter as orp + + def fake_get(url, headers=None, timeout=None, **_kw): + if (headers or {}).get("Authorization") == "Bearer good-key": + return _Resp({"data": [{"id": "google/gemini-image"}]}) + return _Resp({"error": "unauthorized"}, status=401) + + monkeypatch.setattr(requests, "get", fake_get) + orp._CATALOG_CACHE.clear() + try: + assert orp._fetch_image_api_catalog("https://openrouter.ai/api/v1", "bad-key") == frozenset() + assert "google/gemini-image" in orp._fetch_image_api_catalog("https://openrouter.ai/api/v1", "good-key") + finally: + orp._CATALOG_CACHE.clear() + + +def test_memory_provider_skill_prune_only_touches_the_active_home(homes, monkeypatch): + """Pruning under profile B (whose active provider differs) must leave profile A's registered + provider skill in place; A's own later prune still retracts it.""" + import plugins.memory as mem + from hermes_cli.plugins import _reset_plugin_managers_for_tests, get_plugin_manager + + a, b = homes + _reset_plugin_managers_for_tests() + mem._REGISTERED_MEMORY_PROVIDER_SKILLS.clear() + skill_dir = a / "plugins" / "provA" / "skills" / "maint" + skill_dir.mkdir(parents=True) + (skill_dir / "SKILL.md").write_text("---\nname: maint\ndescription: x\n---\nbody\n", encoding="utf-8") + try: + with scoped(a): + mem._ProviderCollector("provA").register_skill("maint", skill_dir) + assert get_plugin_manager().find_plugin_skill("provA:maint") is not None + with scoped(b): + mem._prune_inactive_memory_provider_skills("provB") + with scoped(a): + assert get_plugin_manager().find_plugin_skill("provA:maint") is not None + mem._prune_inactive_memory_provider_skills("provOther") + assert get_plugin_manager().find_plugin_skill("provA:maint") is None + finally: + mem._REGISTERED_MEMORY_PROVIDER_SKILLS.clear() + _reset_plugin_managers_for_tests() + + +def test_openviking_atexit_commits_every_profile_provider(homes): + """Two profiles' providers initialized in one process both get the atexit commit.""" + import plugins.memory.openviking as ov + + a, b = homes + committed: list[object] = [] + providers = [] + try: + for home in (a, b): + with scoped(home): + provider = ov.OpenVikingMemoryProvider() + provider.initialize(session_id=f"s-{home.name}", hermes_home=str(home)) + provider.on_session_end = lambda _msgs, _p=provider: committed.append(_p) + providers.append(provider) + ov._atexit_commit_sessions() + assert committed == providers + finally: + for provider in providers: + with contextlib.suppress(Exception): + provider._release_run_lock() + + +def test_disk_cleanup_protected_cron_paths_follow_the_active_home(homes): + """The protected-path guard must protect the ACTIVE profile's cron dir, not the first one asked.""" + spec = importlib.util.spec_from_file_location( + "disk_cleanup_mux_scope", REPO / "plugins" / "disk-cleanup" / "disk_cleanup.py") + dc = importlib.util.module_from_spec(spec) + spec.loader.exec_module(dc) + a, b = homes + for home in (a, b): + (home / "cron").mkdir() + with scoped(a): + assert dc._is_protected_cron_path(a / "cron") + with scoped(b): + assert dc._is_protected_cron_path(b / "cron") + assert not dc._is_protected_cron_path(a / "cron") + + +def test_yuanbao_active_adapter_resolves_per_profile(homes, monkeypatch): + """Each profile's turn reads back its own adapter; the unscoped slot still serves single-profile.""" + from gateway.config import PlatformConfig + from gateway.platforms.yuanbao import YuanbaoAdapter + + a, b = homes + cfg = PlatformConfig(enabled=True, extra={"app_id": "x", "app_secret": "y"}) + monkeypatch.setattr(YuanbaoAdapter, "_active_instance", None) + monkeypatch.setattr(YuanbaoAdapter, "_active_instances", {}) + with scoped(a): + adapter_a = YuanbaoAdapter(cfg) + YuanbaoAdapter.set_active(adapter_a) + with scoped(b): + adapter_b = YuanbaoAdapter(cfg) + YuanbaoAdapter.set_active(adapter_b) + with scoped(a): + assert YuanbaoAdapter.get_active() is adapter_a + with scoped(b): + assert YuanbaoAdapter.get_active() is adapter_b + assert YuanbaoAdapter.get_active() is None # scoped adapters never claim the unscoped slot + unscoped = YuanbaoAdapter(cfg) + YuanbaoAdapter.set_active(unscoped) + assert YuanbaoAdapter.get_active() is unscoped + + +def test_honcho_loopback_flow_status_is_per_profile(homes, monkeypatch): + """Profile B's connect must not be refused as 'pending' because profile A's flow is running.""" + import plugins.memory.honcho.oauth_flow as flow + + a, b = homes + gate = threading.Event() + started: list[Path] = [] + + def fake_authorize(**kwargs): + started.append(kwargs["config_path"]) + gate.wait(5) + + monkeypatch.setattr(flow, "authorize_via_loopback", fake_authorize) + monkeypatch.setattr(flow, "_status", flow.FlowStatus()) + monkeypatch.setattr(flow, "_flow_thread", None) + for home in (a, b): + (home / "honcho.json").write_text("{}", encoding="utf-8") + try: + with scoped(a): + assert flow.start_loopback_flow_background()["state"] == "pending" + with scoped(b): + assert flow.get_flow_status()["state"] == "idle" + assert flow.start_loopback_flow_background()["state"] == "pending" + deadline = time.monotonic() + 5 + while time.monotonic() < deadline and len(started) < 2: + time.sleep(0.02) + assert sorted(started) == sorted([a / "honcho.json", b / "honcho.json"]) + finally: + gate.set() + getattr(flow, "_flows_by_target", {}).clear()