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.
This commit is contained in:
Teknium
2026-09-12 01:03:46 -07:00
parent 5ff34f565e
commit 208bd0b65a
11 changed files with 544 additions and 94 deletions
+23 -3
View File
@@ -2607,14 +2607,31 @@ class YuanbaoAdapter(BasePlatformAdapter):
MEDIA_MAX_SIZE_MB: int = 50 MEDIA_MAX_SIZE_MB: int = 50
DM_MAX_CHARS = 10000 DM_MAX_CHARS = 10000
_active_instance: ClassVar[Optional["YuanbaoAdapter"]] = None _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 @classmethod
def get_active(cls) -> Optional["YuanbaoAdapter"]: 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 @classmethod
def set_active(cls, adapter: Optional["YuanbaoAdapter"]) -> None: 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: def __init__(self, config: PlatformConfig, **kwargs: Any) -> None:
super().__init__(config, Platform.YUANBAO) super().__init__(config, Platform.YUANBAO)
@@ -2697,7 +2714,10 @@ class YuanbaoAdapter(BasePlatformAdapter):
async def disconnect(self) -> None: async def disconnect(self) -> None:
"""Cancel background tasks and close the WebSocket connection.""" """Cancel background tasks and close the WebSocket connection."""
if YuanbaoAdapter._active_instance is self: 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._running = False
self._mark_disconnected() self._mark_disconnected()
self._release_platform_lock() self._release_platform_lock()
+4 -4
View File
@@ -102,20 +102,20 @@ _NEVER_TRACK_TOP_LEVEL = frozenset({
"patches", "projects", "skins", "themes", "contributors", "patches", "projects", "skins", "themes", "contributors",
"profiles", "backups", "optional-skills"}) "profiles", "backups", "optional-skills"})
@functools.lru_cache(maxsize=1) # built lazily so HERMES_HOME resolves once @functools.lru_cache(maxsize=8) # keyed by home: a multiplexed process serves several profiles
def _protected_cron_paths() -> frozenset: def _protected_cron_paths(home: Path) -> frozenset:
"""Defense-in-depth for quick(): EXACT cron control-plane paths (``cron/``, ``output/`` root, """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). ``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 Never widen to everything under ``cron/output/``: run artifacts there are disposable; only
wholesale deletion of ``output/`` is fatal.""" 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")) 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 # 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. # says. This is a defense-in-depth guard against stale tracked.json entries from before #34840.
def _is_protected_cron_path(p: Path) -> bool: 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: def fmt_size(n: float) -> str:
+11 -6
View File
@@ -62,9 +62,11 @@ _load_image_gen_config = load_image_gen_config
_IMAGE_API_ENV_PREFIX = "OPENROUTER_IMAGE_API_" _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. # 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 _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_TTL_SECONDS = 900.0
_CATALOG_CACHE: Dict[str, Tuple[float, frozenset]] = {} _CATALOG_CACHE: Dict[Tuple[str, Optional[str]], Tuple[float, frozenset]] = {}
_GEMINI_RATIOS = ( _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", "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: 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 """Model ids from ``GET {base_url}/images/models``, cached per (base URL, key). Any failure caches
empty set (→ chat-completions): guessing "images" would 404 a working chat setup.""" an empty set (→ chat-completions): guessing "images" would 404 a working chat setup."""
cached = _CATALOG_CACHE.get(base_url) 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: if cached and (time.monotonic() - cached[0]) < _CATALOG_TTL_SECONDS:
return cached[1] return cached[1]
ids: set = set() 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 except Exception as exc: # noqa: BLE001 - probe must never break generation
logger.debug("image API catalog probe failed for %s: %s", base_url, exc) logger.debug("image API catalog probe failed for %s: %s", base_url, exc)
resolved = frozenset(ids) resolved = frozenset(ids)
_CATALOG_CACHE[base_url] = (time.monotonic(), resolved) _CATALOG_CACHE[cache_key] = (time.monotonic(), resolved)
return resolved return resolved
+33 -7
View File
@@ -43,6 +43,10 @@ _EDIT_FALLBACK_MODEL = "grok-imagine-image-quality"
# Live catalog cache ``(models, fetched_monotonic)``: ``/image-generation-models`` is the source of # 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. # 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 _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_CACHE_TTL = 300.0
_LIVE_TIMEOUT = 10.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("/") 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.""" """``{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() api_key = str(creds.get("api_key") or "").strip()
if not api_key: if not api_key:
raise RuntimeError("no xAI credentials") 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]]: def _live_models() -> Dict[str, Dict[str, Any]]:
"""Cached live catalog (``{}`` when unreachable).""" """Cached live catalog (``{}`` when unreachable)."""
global _LIVE_CACHE 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] return _LIVE_CACHE[0]
from agent.credential_persistence import fingerprint_secret_value
try: 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 except Exception as exc: # noqa: BLE001 - offline/unauth → static fallback
logger.debug("xAI live image model catalog unavailable: %s", exc) logger.debug("xAI live image model catalog unavailable: %s", exc)
live = {} return {}
_LIVE_CACHE = (live, time.monotonic())
return live
def _catalog() -> Dict[str, Dict[str, Any]]: def _catalog() -> Dict[str, Dict[str, Any]]:
+13 -4
View File
@@ -28,7 +28,15 @@ logger = logging.getLogger(__name__)
_MEMORY_PLUGINS_DIR = Path(__file__).parent _MEMORY_PLUGINS_DIR = Path(__file__).parent
ENTRY_POINTS_GROUP = "hermes_agent.memory_providers" 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. # Synthetic parent package so user-installed providers don't collide with bundled ones.
_USER_NAMESPACE = "_hermes_user_memory" _USER_NAMESPACE = "_hermes_user_memory"
@@ -330,7 +338,7 @@ class _ProviderCollector:
registered_path = get_plugin_manager().find_plugin_skill(qualified_name) registered_path = get_plugin_manager().find_plugin_skill(qualified_name)
if registered_path is not None: 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: except Exception as exc:
logger.debug("Memory provider '%s' failed to register skill: %s", self.name, 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 from hermes_cli.plugins import get_plugin_manager
manager = 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: if qualified_name.partition(":")[0] == active_provider:
continue continue
if manager.find_plugin_skill(qualified_name) == registered_path: if manager.find_plugin_skill(qualified_name) == registered_path:
manager.remove_plugin_skill(qualified_name) 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]: def discover_plugin_cli_commands() -> List[dict]:
+11 -6
View File
@@ -91,9 +91,11 @@ def _maybe_upgrade_client() -> None:
pass # packaging not available or other issue — proceed anyway pass # packaging not available or other issue — proceed anyway
# update_mode='append' capability (Hindsight >= 0.5.0), cached per API URL per # update_mode='append' capability (Hindsight >= 0.5.0), cached per (API URL, key fingerprint)
# process so every provider on the same API shares one /version round trip. # per process so every provider on the same API+key shares one /version round trip. A failed probe
_append_capability_cache: Dict[str, bool] = {} # 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() _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: if not api_url:
return False return False
from agent.credential_persistence import fingerprint_secret_value
cache_key = (api_url, fingerprint_secret_value(api_key))
with _append_capability_lock: with _append_capability_lock:
if api_url in _append_capability_cache: if cache_key in _append_capability_cache:
return _append_capability_cache[api_url] return _append_capability_cache[cache_key]
version = _fetch_hindsight_api_version(api_url, api_key) version = _fetch_hindsight_api_version(api_url, api_key)
try: # missing/invalid version -> unsupported try: # missing/invalid version -> unsupported
from packaging.version import Version 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 supported = False
with _append_capability_lock: with _append_capability_lock:
# A concurrent probe may have filled the cache meanwhile; its answer wins. # 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: if supported:
logger.debug("Hindsight API %s version %s supports update_mode='append'", api_url, version) logger.debug("Hindsight API %s version %s supports update_mode='append'", api_url, version)
else: else:
+38 -12
View File
@@ -364,6 +364,25 @@ class FlowStatus:
_status = FlowStatus() _status = FlowStatus()
_status_lock = threading.Lock() _status_lock = threading.Lock()
_flow_thread: threading.Thread | None = None _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]: def _detect_connection() -> tuple[bool, str | None]:
"""Report whether a credential is already stored: 'oauth', 'apikey', or 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 return auth is not None, auth
def get_flow_status() -> dict[str, object]: def get_flow_status() -> dict[str, object]:
status, _thread = _flow_state(_flow_target())
with _status_lock: with _status_lock:
state, detail = _status.state, _status.detail state, detail = status.state, status.detail
connected, auth = _detect_connection() connected, auth = _detect_connection()
return {"state": state, "detail": detail, "connected": connected, "auth": auth} 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: with _status_lock:
_status.state, _status.detail = state, detail status.state, status.detail = state, detail
def start_loopback_flow_background( def start_loopback_flow_background(
*, config_path: Path | None = None, host: str | None = None, source: str = "hermes-desktop", *, 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.""" Idempotent while pending, so a double-click can't open two tabs / bind :8765 twice."""
global _flow_thread global _flow_thread
# Resolve under the caller's profile scope NOW — a context-local HERMES_HOME override can't reach the worker. # 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() target = _flow_target()
host = host or resolve_active_host() 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: with _status_lock:
if _status.state == "pending" and _flow_thread and _flow_thread.is_alive(): if status.state == "pending" and thread and thread.is_alive():
return {"state": _status.state, "detail": _status.detail} return {"state": status.state, "detail": status.detail}
_status.state, _status.detail = "pending", "waiting for browser consent" status.state, status.detail = "pending", "waiting for browser consent"
def _run() -> None: def _run() -> None:
try: try:
authorize_via_loopback(config_path=config_path, host=host, source=source, timeout=timeout) 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: except Exception as exc:
logger.warning("Honcho OAuth loopback flow failed: %s", 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) thread = threading.Thread(target=_run, name="honcho-oauth-loopback", daemon=True)
_flow_thread.start() if target is None:
_flow_thread = thread
else:
_flows_by_target[target] = (status, thread)
thread.start()
return get_flow_status() return get_flow_status()
+18 -19
View File
@@ -198,23 +198,23 @@ def _preview(value: Any, limit: int = 160) -> str:
# atexit safety net: commit pending sessions even if shutdown_memory_provider # atexit safety net: commit pending sessions even if shutdown_memory_provider
# never runs (gateway crash, exception in the session expiry watcher, ...). # 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(): def _atexit_commit_sessions():
global _last_active_provider providers = list(_active_providers_by_home.values())
provider = _last_active_provider _active_providers_by_home.clear()
if provider is None: for provider in providers:
return try:
_last_active_provider = None with suppress(Exception): # best-effort at shutdown time
try: provider.on_session_end([])
with suppress(Exception): # best-effort at shutdown time finally:
provider.on_session_end([]) # ``finally`` (as on main): the run lock is released even when on_session_end
finally: # dies of a BaseException (KeyboardInterrupt during atexit).
# ``finally`` (as on main): the run lock is released even when on_session_end with suppress(Exception):
# dies of a BaseException (KeyboardInterrupt during atexit). provider._release_run_lock()
with suppress(Exception):
provider._release_run_lock()
atexit.register(_atexit_commit_sessions) atexit.register(_atexit_commit_sessions)
@@ -1436,8 +1436,7 @@ class OpenVikingMemoryProvider(MemoryProvider):
self._conn_snapshot = self._settings_tuple() self._conn_snapshot = self._settings_tuple()
self._recover_pending_sessions() self._recover_pending_sessions()
global _last_active_provider # atexit safety net _active_providers_by_home[self._hermes_home] = self # atexit safety net
_last_active_provider = self
def _ensure_client(self) -> Optional["_VikingClient"]: def _ensure_client(self) -> Optional["_VikingClient"]:
"""Active client, rebuilt if the resolved config changed. """Active client, rebuilt if the resolved config changed.
@@ -2488,9 +2487,9 @@ class OpenVikingMemoryProvider(MemoryProvider):
for t in workers: for t in workers:
if t.is_alive(): if t.is_alive():
t.join(timeout=5.0) t.join(timeout=5.0)
global _last_active_provider # clear so atexit doesn't double-commit # Clear so atexit doesn't double-commit.
if _last_active_provider is self: if _active_providers_by_home.get(self._hermes_home) is self:
_last_active_provider = None del _active_providers_by_home[self._hermes_home]
self._release_run_lock() self._release_run_lock()
@staticmethod @staticmethod
+61 -17
View File
@@ -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``. request hot path) and fed to the codex transport's clamp via ``supported_reasoning_efforts``.
""" """
import contextvars
import json import json
import logging import logging
import os import os
import sys
import threading import threading
import time import time
from pathlib import Path from pathlib import Path
@@ -32,13 +34,49 @@ _efforts_lock = threading.Lock()
_warm_started = False _warm_started = False
_disk_checked = False _disk_checked = False
# A stale verdict beats no verdict: a past-TTL mirror is still served while a
# background refresh runs. class _CacheState:
_DISK_TTL_SECONDS = 24 * 60 * 60 """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: 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: 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]]]: def _seed_efforts(items: Any) -> Optional[dict[str, list[str]]]:
"""Seed memory + disk caches from a ``/v1/models`` payload.""" """Seed memory + disk caches from a ``/v1/models`` payload."""
global _efforts_cache
parsed = _parse_efforts(items) parsed = _parse_efforts(items)
if parsed is not None: if parsed is not None:
state = _state()
with _efforts_lock: with _efforts_lock:
_efforts_cache = parsed state._efforts_cache = parsed
_save_disk(parsed) _save_disk(parsed)
return 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]]]: def _efforts_cache_only() -> Optional[dict[str, list[str]]]:
"""Memory, else the disk mirror (checked once per process). Never HTTP (hot-path safe).""" """Memory, else the disk mirror (checked once per home). Never HTTP (hot-path safe)."""
global _efforts_cache, _disk_checked state = _state()
with _efforts_lock: with _efforts_lock:
cached = _efforts_cache cached = state._efforts_cache
if cached is not None or _disk_checked: if cached is not None or state._disk_checked:
return cached return cached
_disk_checked = True state._disk_checked = True
parsed, age = _load_disk() parsed, age = _load_disk()
if parsed is None: if parsed is None:
return None return None
with _efforts_lock: 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: if age >= _DISK_TTL_SECONDS:
_warm_efforts_async() _warm_efforts_async()
return cached return cached
def _warm_efforts_async() -> None: 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) 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). and without a key (it would 401; the first authenticated fetch_models() seeds).
""" """
global _warm_started
if os.environ.get("PYTEST_CURRENT_TEST"): if os.environ.get("PYTEST_CURRENT_TEST"):
return return
state = _state()
with _efforts_lock: with _efforts_lock:
if _warm_started: if state._warm_started:
return return
_warm_started = True state._warm_started = True
if not _resolve_api_key(): if not _resolve_api_key():
return return
@@ -211,7 +251,11 @@ def _warm_efforts_async() -> None:
if items is not None: if items is not None:
_seed_efforts(items) _seed_efforts(items)
try: 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: except Exception as exc:
logger.debug("router: caps warmer failed to start: %s", exc) logger.debug("router: caps warmer failed to start: %s", exc)
+60 -16
View File
@@ -52,6 +52,10 @@ _TRACE_STATE: Dict[str, TraceState] = {}
# Bounds the leak, not concurrency. # Bounds the leak, not concurrency.
_MAX_TRACE_STATE = 256 _MAX_TRACE_STATE = 256
_LANGFUSE_CLIENT = None _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 # 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. # first client build so racing callers can't each construct a client.
_LANGFUSE_CLIENT_LOCK = threading.Lock() _LANGFUSE_CLIENT_LOCK = threading.Lock()
@@ -83,6 +87,19 @@ def _env(name: str, default: str = "") -> str:
return os.environ.get(name, default).strip() 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: def _debug(message: str) -> None:
if _env("HERMES_LANGFUSE_DEBUG").lower() in {"1", "true", "yes", "on"}: if _env("HERMES_LANGFUSE_DEBUG").lower() in {"1", "true", "yes", "on"}:
logger.info("Langfuse tracing: %s", message) 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)" 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]: def _get_langfuse() -> Optional[Langfuse]:
"""Cached Langfuse client, or ``None`` if the SDK/credentials are unavailable. """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 The first build is serialized so racing callers can't each construct a client
and leak the loser's HTTP connection + flush thread.""" and leak the loser's HTTP connection + flush thread."""
global _LANGFUSE_CLIENT
# Fast path — already settled (success or _INIT_FAILED) needs no lock; # Fast path — already settled (success or _INIT_FAILED) needs no lock;
# re-check under it since a racing thread may have finished init. # 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: with _LANGFUSE_CLIENT_LOCK:
if _LANGFUSE_CLIENT is None: settled = _settled_client()
client = _build_client() if settled is None:
_LANGFUSE_CLIENT = _INIT_FAILED if client is None else client settled = _settle_client()
if client is not None: return None if settled is _INIT_FAILED else settled
# 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
def _build_client() -> Optional[Langfuse]: def _build_client() -> Optional[Langfuse]:
@@ -211,7 +252,7 @@ def _build_client() -> Optional[Langfuse]:
) )
return None 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): if not (public_key and secret_key):
return None return None
@@ -234,10 +275,10 @@ def _build_client() -> Optional[Langfuse]:
kwargs: Dict[str, Any] = {"public_key": public_key, "secret_key": secret_key} 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", ""), for key, name, default in (("base_url", "BASE_URL", "https://cloud.langfuse.com"), ("environment", "ENV", ""),
("release", "RELEASE", "")): ("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: if value:
kwargs[key] = value kwargs[key] = value
sample_rate = _env("HERMES_LANGFUSE_SAMPLE_RATE") sample_rate = _secret("HERMES_LANGFUSE_SAMPLE_RATE")
if sample_rate: if sample_rate:
try: try:
kwargs["sample_rate"] = float(sample_rate) kwargs["sample_rate"] = float(sample_rate)
@@ -596,7 +637,10 @@ def _finalize_all_traces() -> None:
_end_children(state, include_subagents=True) _end_children(state, include_subagents=True)
_end_root(state, f"atexit finalize for {key}") _end_root(state, f"atexit finalize for {key}")
if states: 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: 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 tool-only or empty final response never reaches ``_finish_trace``; its root
would dangle until eviction and queued events could be lost on exit.""" 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. # 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"): if client is None or client is _INIT_FAILED or not hasattr(client, "flush"):
return return
@@ -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()