diff --git a/agent/auxiliary_client.py b/agent/auxiliary_client.py index c92af45a7e..18e5d123db 100644 --- a/agent/auxiliary_client.py +++ b/agent/auxiliary_client.py @@ -5861,8 +5861,10 @@ def _get_task_extra_body(task: str) -> Dict[str, Any]: # During provider incidents each call also retries / fans out across the fallback chain, multiplying request # volume on already-degraded endpoints. A per-task semaphore caps in-flight calls so retry amplification # stays bounded. See #23324. -_aux_sync_semaphores: Dict[str, Tuple[int, threading.BoundedSemaphore]] = {} -_aux_async_semaphores: Dict[Tuple[str, int], Tuple[int, Any]] = {} +# Keyed by profile home as well: the limit is the profile's ``auxiliary..max_concurrency``, and two +# multiplexed profiles with different limits would otherwise rebuild (and reset) one shared semaphore. +_aux_sync_semaphores: Dict[Tuple[str, str], Tuple[int, threading.BoundedSemaphore]] = {} +_aux_async_semaphores: Dict[Tuple[str, str, int], Tuple[int, Any]] = {} _aux_sem_lock = threading.Lock() @@ -5890,7 +5892,10 @@ def _cached_semaphore(store: dict, key: Any, limit: int, factory: Callable[[int] def _acquire_sync_aux_semaphore(task: Optional[str]) -> Optional[threading.BoundedSemaphore]: """Get a per-task sync semaphore, rebuilding it after a config change.""" limit = _get_task_max_concurrency(task) - return None if limit is None else _cached_semaphore(_aux_sync_semaphores, task, limit, threading.BoundedSemaphore) + if limit is None: + return None + from hermes_constants import hermes_home_key + return _cached_semaphore(_aux_sync_semaphores, (hermes_home_key(), task), limit, threading.BoundedSemaphore) def _acquire_async_aux_semaphore(task: Optional[str]): @@ -5903,7 +5908,8 @@ def _acquire_async_aux_semaphore(task: Optional[str]): loop = asyncio.get_running_loop() except RuntimeError: return None - return _cached_semaphore(_aux_async_semaphores, (task, id(loop)), limit, asyncio.Semaphore) + from hermes_constants import hermes_home_key + return _cached_semaphore(_aux_async_semaphores, (hermes_home_key(), task, id(loop)), limit, asyncio.Semaphore) def _reset_aux_semaphores() -> None: diff --git a/agent/image_token_cost.py b/agent/image_token_cost.py index 1fc5e570ee..78219bf93d 100644 --- a/agent/image_token_cost.py +++ b/agent/image_token_cost.py @@ -28,6 +28,9 @@ _EMA_ALPHA = 0.5 _image_cost_var: ContextVar[Optional[int]] = ContextVar("hermes_image_token_cost", default=None) _LEARNED: Dict[str, int] = {} _LOADED = False +# Routed profiles (multiplexed gateway) keep their own table, loaded from THEIR cache file: the +# module slot above is the launch profile's and would otherwise be persisted into every home. +_LEARNED_BY_HOME: Dict[str, Dict[str, int]] = {} def _cache_path(): @@ -42,22 +45,33 @@ def _key(model: Any, base_url: Any) -> str: return f"{model or ''}@{base_url_hostname(base_url or '') or ''}" -def _load() -> None: - global _LOADED - if _LOADED: - return - _LOADED = True +def _read_cache() -> Dict[str, int]: from agent.model_metadata import _load_json_dict - for k, v in _load_json_dict(_cache_path()).items(): - if isinstance(v, int) and _MIN_PLAUSIBLE <= v <= _MAX_PLAUSIBLE: - _LEARNED[k] = v + return {k: v for k, v in _load_json_dict(_cache_path()).items() + if isinstance(v, int) and _MIN_PLAUSIBLE <= v <= _MAX_PLAUSIBLE} + + +def _table() -> Dict[str, int]: + """The active profile's learned table, loaded lazily from its cache file.""" + global _LOADED + from hermes_constants import get_hermes_home_override, hermes_home_key + + if get_hermes_home_override() is None: + if not _LOADED: + _LOADED = True + _LEARNED.update(_read_cache()) + return _LEARNED + home_key = hermes_home_key() + table = _LEARNED_BY_HOME.get(home_key) + if table is None: + table = _LEARNED_BY_HOME[home_key] = _read_cache() + return table def learned_image_token_cost(model: Any, base_url: Any) -> int: """Learned per-image cost for ``model@host``, else the flat default.""" - _load() - return _LEARNED.get(_key(model, base_url), DEFAULT_IMAGE_TOKEN_COST) + return _table().get(_key(model, base_url), DEFAULT_IMAGE_TOKEN_COST) def current_image_token_cost() -> int: @@ -117,15 +131,15 @@ def calibrate_from_usage(agent: Any, messages: List[Dict[str, Any]], prompt_toke if not _MIN_PLAUSIBLE <= per_image <= _MAX_PLAUSIBLE: return None key = _key(getattr(agent, "model", None), getattr(agent, "base_url", None)) - _load() - prior = _LEARNED.get(key) + table = _table() + prior = table.get(key) learned = per_image if prior is None else int(prior + _EMA_ALPHA * (per_image - prior)) - _LEARNED[key] = learned + table[key] = learned _image_cost_var.set(learned) try: from utils import atomic_json_write - atomic_json_write(_cache_path(), dict(_LEARNED), indent=0, separators=(",", ":")) + atomic_json_write(_cache_path(), dict(table), indent=0, separators=(",", ":")) except Exception: logger.debug("image token cost persist failed", exc_info=True) logger.info( diff --git a/agent/model_metadata.py b/agent/model_metadata.py index a36092ce8d..402b400e70 100644 --- a/agent/model_metadata.py +++ b/agent/model_metadata.py @@ -105,8 +105,10 @@ def _strip_provider_prefix(model: str) -> str: _model_metadata_cache: Dict[str, Dict[str, Any]] = {} _model_metadata_cache_time: float = 0 _MODEL_CACHE_TTL = 3600 -_endpoint_model_metadata_cache: Dict[str, Dict[str, Dict[str, Any]]] = {} -_endpoint_model_metadata_cache_time: Dict[str, float] = {} +# In-memory memo keyed by (base_url, api-key fingerprint): per-key gateways return a per-key catalog, and +# in a multiplexed process two profiles may share a URL with different keys. The disk memo stays per URL. +_endpoint_model_metadata_cache: Dict[Tuple[str, str], Dict[str, Dict[str, Any]]] = {} +_endpoint_model_metadata_cache_time: Dict[Tuple[str, str], float] = {} _ENDPOINT_MODEL_CACHE_TTL = 300 # Server-type verdicts (server_type, monotonic_ts): positive ones live an hour so a # server swap on the same port is re-detected; None gets the short TTL so a @@ -940,9 +942,15 @@ def _apply_llamacpp_props(cache: Dict[str, Dict[str, Any]], request_candidate: s cache[child_id]["context_length"] = child_ctx -def _remember_endpoint_models(normalized: str, cache: Dict[str, Dict[str, Any]]) -> Dict[str, Dict[str, Any]]: - _endpoint_model_metadata_cache[normalized] = cache - _endpoint_model_metadata_cache_time[normalized] = time.time() +def _endpoint_memo_key(normalized: str, api_key: object) -> Tuple[str, str]: + from agent.credential_persistence import fingerprint_secret_value + # Callable (minted) keys are not fingerprinted here: doing so would mint on every cache hit. + return normalized, (fingerprint_secret_value(api_key) or "") if isinstance(api_key, str) else "" + + +def _remember_endpoint_models(memo_key: Tuple[str, str], cache: Dict[str, Dict[str, Any]]) -> Dict[str, Dict[str, Any]]: + _endpoint_model_metadata_cache[memo_key] = cache + _endpoint_model_metadata_cache_time[memo_key] = time.time() return cache @@ -962,13 +970,14 @@ def fetch_endpoint_model_metadata(base_url: str, api_key: str = "", force_refres return {} _ensure_requests() local = is_local_endpoint(normalized) + memo_key = _endpoint_memo_key(normalized, api_key) if not force_refresh: - cached = _endpoint_model_metadata_cache.get(normalized) - if cached is not None and (time.time() - _endpoint_model_metadata_cache_time.get(normalized, 0)) < _ENDPOINT_MODEL_CACHE_TTL: + cached = _endpoint_model_metadata_cache.get(memo_key) + if cached is not None and (time.time() - _endpoint_model_metadata_cache_time.get(memo_key, 0)) < _ENDPOINT_MODEL_CACHE_TTL: return cached memo = _endpoint_disk_cache_get(normalized) if not local else None if memo is not None: - return _remember_endpoint_models(normalized, memo) + return _remember_endpoint_models(memo_key, memo) # Blackholed: return empty WITHOUT caching so it is retried once the entry expires. if _endpoint_blackholed(normalized): return {} @@ -980,7 +989,7 @@ def fetch_endpoint_model_metadata(base_url: str, api_key: str = "", force_refres if local: try: if detect_local_server_type(normalized, api_key=api_key) == "lm-studio": - return _remember_endpoint_models(normalized, _lmstudio_native_models(normalized, headers)) + return _remember_endpoint_models(memo_key, _lmstudio_native_models(normalized, headers)) except Exception as exc: last_error = exc _note_if_connect_timeout(exc, normalized) @@ -1005,7 +1014,7 @@ def fetch_endpoint_model_metadata(base_url: str, api_key: str = "", force_refres _apply_llamacpp_props(cache, request_candidate, headers, verify) if cache and not local: _endpoint_disk_cache_put(normalized, cache) - return _remember_endpoint_models(normalized, cache) + return _remember_endpoint_models(memo_key, cache) except Exception as exc: last_error = exc _note_if_connect_timeout(exc, normalized) @@ -1014,7 +1023,7 @@ def fetch_endpoint_model_metadata(base_url: str, api_key: str = "", force_refres response.close() if last_error: logger.debug("Failed to fetch model metadata from %s/models: %s", normalized, last_error) - return _remember_endpoint_models(normalized, {}) + return _remember_endpoint_models(memo_key, {}) def _resolve_endpoint_context_length(model: str, base_url: str, api_key: str = "") -> Optional[int]: diff --git a/agent/prompt_builder.py b/agent/prompt_builder.py index 77db1c348a..7913cf1025 100644 --- a/agent/prompt_builder.py +++ b/agent/prompt_builder.py @@ -798,8 +798,10 @@ _BACKEND_FALLBACK_DESCRIPTIONS: dict[str, str] = { "ssh": "a remote host reached over SSH (likely Linux)", } -# Per-process probe cache keyed by (env_type, cwd_hint) so a mid-process backend switch rebuilds. -_BACKEND_PROBE_CACHE: dict[tuple[str, str], str] = {} +# Per-process probe cache keyed by (home key, env_type, cwd_hint) so a mid-process backend switch +# rebuilds; the home key because the probe runs against the profile's own terminal.* backend +# (docker image / ssh host) and one multiplexed process serves several profiles. +_BACKEND_PROBE_CACHE: dict[tuple[str, str, str], str] = {} def _plugin_backend_attr(backend: str, attr: str, default=None): @@ -918,7 +920,8 @@ def _format_backend_probe(output: str) -> str: def _probe_remote_backend(env_type: str) -> str | None: """Describe the active non-local backend via a live probe; None if it failed (cached, failures included).""" - cache_key = (env_type, _tenv_read("TERMINAL_CWD", "")) + from hermes_constants import hermes_home_key + cache_key = (hermes_home_key(), env_type, _tenv_read("TERMINAL_CWD", "")) formatted = _BACKEND_PROBE_CACHE.get(cache_key) if formatted is None: formatted = "" diff --git a/tests/tools/test_multiplex_tool_cache_scope.py b/tests/tools/test_multiplex_tool_cache_scope.py new file mode 100644 index 0000000000..7eb0d3691e --- /dev/null +++ b/tests/tools/test_multiplex_tool_cache_scope.py @@ -0,0 +1,187 @@ +"""Multiplexed gateway: module-level caches in tools/ and agent/ must not hand profile A's +config/.env-derived value to profile B. Real temp homes, real config.yaml/.env, real modules; +only HTTP transports are stubbed. +""" +from __future__ import annotations + +import json +import threading +from pathlib import Path + +import pytest +import yaml + +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 + + +def _make_home(root: Path, cfg: dict, env: str = "") -> Path: + root.mkdir(parents=True, exist_ok=True) + (root / "config.yaml").write_text(yaml.safe_dump(cfg), encoding="utf-8") + (root / ".env").write_text(env, encoding="utf-8") + (root / "cache").mkdir(exist_ok=True) + return root + + +class _scoped: + def __init__(self, home: Path): + self.home = home + + def __enter__(self): + self._t1 = set_hermes_home_override(str(self.home)) + self._t2 = set_secret_scope(build_profile_secret_scope(self.home)) + + def __exit__(self, *_): + reset_secret_scope(self._t2) + reset_hermes_home_override(self._t1) + + +@pytest.fixture +def two_homes(tmp_path, monkeypatch): + a = _make_home(tmp_path / "A", {}, "CAMOFOX_URL=http://camofox-a:9377\n") + b = _make_home(tmp_path / "A" / "profiles" / "B", {}, "CAMOFOX_URL=http://camofox-b:9377\n") + monkeypatch.setenv("HERMES_HOME", str(a)) + monkeypatch.delenv("CAMOFOX_URL", raising=False) + return a, b + + +def test_camofox_vnc_memo_is_keyed_by_the_profiles_server_url(two_homes, monkeypatch): + """The one-shot VNC probe must not answer B (own CAMOFOX_URL) with A's server address.""" + import tools.browser_camofox as cam + + class _Resp: + status_code = 200 + + def __init__(self, url): + self._port = 6001 if "camofox-a" in url else 6002 + + def json(self): + return {"ok": True, "vncPort": self._port} + + monkeypatch.setattr(cam.requests, "get", lambda url, *a, **k: _Resp(url)) + a, b = two_homes + with _scoped(a): + assert cam.check_camofox_available() is True + assert cam.get_vnc_url() == "http://camofox-a:6001" + with _scoped(b): + assert cam.get_vnc_url() == "http://camofox-b:6002" + with _scoped(a): + assert cam.get_vnc_url() == "http://camofox-a:6001" + + +def test_home_keyed_caches_serve_each_profile_its_own_config(tmp_path, monkeypatch): + """One mechanism (dict keyed by home / override bypass) across the sites that read per-profile + config or per-home files: aux-vision routing, tirith binary path, learned image cost table, aux + semaphore, MCP lock.""" + bin_a, bin_b = tmp_path / "binA" / "tirith", tmp_path / "binB" / "tirith" + for p in (bin_a, bin_b): + p.parent.mkdir() + p.write_text("#!/bin/sh\nexit 0\n", encoding="utf-8") + p.chmod(0o755) + main = {"model": {"provider": "openai", "model": "gpt-4o"}} + a = _make_home(tmp_path / "A", {**main, "security": {"tirith_path": str(bin_a)}, + "auxiliary": {"vision": {"provider": "auto"}, "summary": {"max_concurrency": 2}}}) + b = _make_home(tmp_path / "A" / "profiles" / "B", {**main, "security": {"tirith_path": str(bin_b)}, + "auxiliary": {"vision": {"provider": "openai", "model": "gpt-4o-mini"}, + "summary": {"max_concurrency": 7}}}) + monkeypatch.setenv("HERMES_HOME", str(a)) + (a / "cache" / "image_token_costs.json").write_text(json.dumps({"m@gw.example": 1000}), encoding="utf-8") + (b / "cache" / "image_token_costs.json").write_text(json.dumps({"m@gw.example": 3000}), encoding="utf-8") + + import agent.auxiliary_client as ac + import agent.image_token_cost as itc + import tools.computer_use.tool as cu + import tools.tirith_security as tir + from tools import mcp_tool_loop + + monkeypatch.setattr(tir, "_resolved_path", None) + monkeypatch.setattr(tir, "_resolved_path_by_home", {}) + ac._reset_aux_semaphores() + cu._AUX_VISION_ROUTE_CACHE.clear() + + with _scoped(a): + assert cu._should_route_through_aux_vision() is False # no explicit aux vision: native path + assert tir._resolve_tirith_path(tir._load_security_config()["tirith_path"]) == str(bin_a) + assert itc.learned_image_token_cost("m", "http://gw.example/v1") == 1000 + sem_a = ac._acquire_sync_aux_semaphore("summary") + sem_a.acquire() + cookie = mcp_tool_loop._try_acquire_mcp_discovery_lock() + lock_a = cookie._fh.name + cookie.release() + with _scoped(b): + assert cu._should_route_through_aux_vision() is True # B named a dedicated vision model + assert tir._resolve_tirith_path(tir._load_security_config()["tirith_path"]) == str(bin_b) + assert itc.learned_image_token_cost("m", "http://gw.example/v1") == 3000 + sem_b = ac._acquire_sync_aux_semaphore("summary") + cookie = mcp_tool_loop._try_acquire_mcp_discovery_lock() + lock_b = cookie._fh.name + cookie.release() + assert Path(lock_a).parent == a and Path(lock_b).parent == b + assert sem_b is not sem_a + with _scoped(a): + # B's differently-sized lookup must not have rebuilt the semaphore A is holding. + assert ac._acquire_sync_aux_semaphore("summary") is sem_a + sem_a.release() + + +def test_debounced_sync_push_fires_in_the_scheduling_profiles_context(two_homes, monkeypatch): + """Timer threads start with empty ContextVars: the push must run under the writing profile's + home, and B's write must not cancel A's pending push.""" + import tools.skill_manager_tool as smt + import tools.skill_usage as su + import tools.skills_sync_client as ssc + from hermes_constants import get_hermes_home + + a, b = two_homes + fired: dict[str, str] = {} + both = threading.Event() + + def fake_push(*, message=""): + fired[message] = str(get_hermes_home()) + if len(fired) == 2: + both.set() + + monkeypatch.setattr(su, "is_sync_enabled", lambda name: True) + monkeypatch.setattr(ssc, "maybe_push_skills", fake_push) + monkeypatch.setattr(smt, "_SYNC_PUSH_DEBOUNCE_S", 0.05) + monkeypatch.setattr(smt, "_sync_push_timers", {}) + with _scoped(a): + smt._maybe_debounced_sync_push("skill-a") + with _scoped(b): + smt._maybe_debounced_sync_push("skill-b") + assert both.wait(5), fired + assert fired == {"sync: skill-a": str(a), "sync: skill-b": str(b)} + + +def test_endpoint_model_catalog_memo_is_keyed_by_credential(two_homes, monkeypatch): + """Two profiles, same base_url, different api_key: a per-key gateway's catalog fetched with A's + key must not be served to B from the in-memory memo (the disk memo already lives per home).""" + import agent.model_metadata as mm + + a, b = two_homes + mm._endpoint_model_metadata_cache.clear() + mm._endpoint_model_metadata_cache_time.clear() + mm._ensure_requests() + + class _Resp: + status_code, ok = 200, True + + def __init__(self, headers): + self._who = headers.get("Authorization", "").rsplit("-", 1)[-1] + + def raise_for_status(self): + pass + + def json(self): + return {"data": [{"id": f"model-for-{self._who}", "context_length": 1}]} + + def close(self): + pass + + monkeypatch.setattr(mm.requests, "get", lambda url, headers=None, **k: _Resp(headers or {})) + with _scoped(a): + assert set(mm.fetch_endpoint_model_metadata("http://gw.example/v1", api_key="key-A")) == {"model-for-A"} + with _scoped(b): + assert set(mm.fetch_endpoint_model_metadata("http://gw.example/v1", api_key="key-B")) == {"model-for-B"} + with _scoped(a): + assert set(mm.fetch_endpoint_model_metadata("http://gw.example/v1", api_key="key-A")) == {"model-for-A"} diff --git a/tools/browser_camofox.py b/tools/browser_camofox.py index 08917a9d64..8618a595c6 100644 --- a/tools/browser_camofox.py +++ b/tools/browser_camofox.py @@ -25,7 +25,7 @@ import requests from agent.secret_scope import get_secret from hermes_cli.config import cfg_get, load_config, read_raw_config -from hermes_constants import hermes_home_key +from hermes_constants import get_hermes_home_override, hermes_home_key from tools.browser_camofox_state import get_camofox_identity from tools.registry import tool_error @@ -37,6 +37,9 @@ _DEFAULT_TIMEOUT = 30 # fallback when config is unreadable _NO_SESSION_ERROR = "No browser session. Call browser_navigate first." _vnc_url: Optional[str] = None # cached from /health response _vnc_url_checked = False # only probe once per process +# Routed profiles (multiplexed gateway) each point CAMOFOX_URL at their own server, so the one-shot +# slot above would hand the launch profile's VNC address to every other profile: memo per server URL. +_vnc_url_by_camofox_url: Dict[str, Optional[str]] = {} # browser.command_timeout, resolved lazily like browser_tool; keyed by profile home because the # multiplexed gateway serves every profile from one process. _cached_cmd_timeout: Optional[Dict[str, int]] = None @@ -107,6 +110,16 @@ def is_camofox_mode() -> bool: return bool(get_camofox_url()) +def _vnc_url_from_health(url: str, resp: Any) -> Optional[str]: + try: + vnc_port = resp.json().get("vncPort") + if isinstance(vnc_port, int) and 1 <= vnc_port <= 65535: + return f"http://{urlsplit(url).hostname or 'localhost'}:{vnc_port}" + except (ValueError, KeyError): + pass + return None + + def check_camofox_available() -> bool: """Verify the Camofox server is reachable (and cache its VNC URL once).""" global _vnc_url, _vnc_url_checked @@ -117,19 +130,23 @@ def check_camofox_available() -> bool: resp = requests.get(f"{url}/health", timeout=5) except Exception: return False - if resp.status_code == 200 and not _vnc_url_checked: - try: - vnc_port = resp.json().get("vncPort") - if isinstance(vnc_port, int) and 1 <= vnc_port <= 65535: - _vnc_url = f"http://{urlsplit(url).hostname or 'localhost'}:{vnc_port}" - except (ValueError, KeyError): - pass - _vnc_url_checked = True + if resp.status_code == 200: + if get_hermes_home_override() is not None: + if url not in _vnc_url_by_camofox_url: + _vnc_url_by_camofox_url[url] = _vnc_url_from_health(url, resp) + elif not _vnc_url_checked: + _vnc_url = _vnc_url_from_health(url, resp) or _vnc_url + _vnc_url_checked = True return resp.status_code == 200 def get_vnc_url() -> Optional[str]: """Return the VNC URL if the Camofox server exposes one, or None.""" + if get_hermes_home_override() is not None: + url = get_camofox_url() + if url not in _vnc_url_by_camofox_url: + check_camofox_available() + return _vnc_url_by_camofox_url.get(url) if not _vnc_url_checked: check_camofox_available() return _vnc_url diff --git a/tools/computer_use/tool.py b/tools/computer_use/tool.py index 953bfda4b9..afb000913f 100644 --- a/tools/computer_use/tool.py +++ b/tools/computer_use/tool.py @@ -79,7 +79,9 @@ _backend: Optional[ComputerUseBackend] = None # backward-compatible empty-sessi _backends: Dict[str, ComputerUseBackend] = {} _backend_call_locks: Dict[str, threading.RLock] = {} _backend_permission_modes: Dict[str, str] = {} -_AUX_VISION_ROUTE_CACHE: Dict[Tuple[str, str], bool] = {} # process-scoped: (provider, model) → bool +# (home key, provider, model) → bool. The decision reads the active profile's config (auxiliary.vision +# override, declared supports_vision), so a multiplexed process must not serve profile A's verdict to B. +_AUX_VISION_ROUTE_CACHE: Dict[Tuple[str, str, str], bool] = {} # Approval state keyed by session_id so a gateway serving concurrent sessions can't leak one run's # "always approve" into another; callers without a session_id share "". # Falls back to a shared "" bucket for callers that don't pass a session_id (e.g. the classic single-run @@ -675,10 +677,11 @@ def _should_route_through_aux_vision() -> bool: try: from agent.auxiliary_client import _read_main_model, _read_main_provider from hermes_cli.config import load_config + from hermes_constants import hermes_home_key from tools.computer_use.vision_routing import should_route_capture_to_aux_vision stage = "config read" provider, model = _read_main_provider() or "", _read_main_model() or "" - if (cached := _AUX_VISION_ROUTE_CACHE.get(key := (str(provider), str(model)))) is not None: + if (cached := _AUX_VISION_ROUTE_CACHE.get(key := (hermes_home_key(), str(provider), str(model)))) is not None: return cached stage = "decision" _AUX_VISION_ROUTE_CACHE[key] = decision = bool(should_route_capture_to_aux_vision(provider, model, load_config())) diff --git a/tools/mcp_tool_loop.py b/tools/mcp_tool_loop.py index d2154d8fff..831ea81bb6 100644 --- a/tools/mcp_tool_loop.py +++ b/tools/mcp_tool_loop.py @@ -67,12 +67,18 @@ def _try_acquire_mcp_discovery_lock() -> Any: """``_LockCookie`` (acquired), ``None`` (held by another process) or ``_LOCK_UNAVAILABLE`` (locking broken: run discovery unguarded).""" # The cached path lives on the ORIGIN module (tests reset ``tools.mcp_tool._MCP_DISCOVERY_LOCK_PATH``). + # A routed profile (multiplexed gateway) locks under ITS home: the launch profile's lock file + # would serialize discovery across profiles and never coordinate with B's own single-profile processes. from tools import mcp_tool as _origin try: - from hermes_constants import get_hermes_home - if _origin._MCP_DISCOVERY_LOCK_PATH is None: - _origin._MCP_DISCOVERY_LOCK_PATH = str(get_hermes_home() / ".mcp-discovery.lock") - fh = open(_origin._MCP_DISCOVERY_LOCK_PATH, "w", encoding="utf-8") + from hermes_constants import get_hermes_home, get_hermes_home_override + if get_hermes_home_override() is not None: + lock_path = str(get_hermes_home() / ".mcp-discovery.lock") + else: + if _origin._MCP_DISCOVERY_LOCK_PATH is None: + _origin._MCP_DISCOVERY_LOCK_PATH = str(get_hermes_home() / ".mcp-discovery.lock") + lock_path = _origin._MCP_DISCOVERY_LOCK_PATH + fh = open(lock_path, "w", encoding="utf-8") except Exception: return _core._LOCK_UNAVAILABLE try: diff --git a/tools/skill_manager_tool.py b/tools/skill_manager_tool.py index c69bdeca3e..00a7ce8c0b 100644 --- a/tools/skill_manager_tool.py +++ b/tools/skill_manager_tool.py @@ -632,7 +632,8 @@ def apply_skill_pending(payload: Dict[str, Any]) -> str: # Sync push debounce: a burst of skill_manage writes collapses into one push on a daemon timer. -_sync_push_timer = None +# One timer per profile home: in a multiplexed process B's write must not cancel A's pending push. +_sync_push_timers: Dict[str, threading.Timer] = {} _sync_push_lock = threading.Lock() _SYNC_PUSH_DEBOUNCE_S = 5.0 @@ -640,23 +641,29 @@ _SYNC_PUSH_DEBOUNCE_S = 5.0 def _maybe_debounced_sync_push(skill_name: str) -> None: """Debounced best-effort sync push after a skill write; never blocks the caller. Skills not opted into sync do nothing (no auth/network); ``maybe_push_skills`` enforces the access gate.""" - global _sync_push_timer try: from tools.skill_usage import is_sync_enabled if not is_sync_enabled(skill_name): return except Exception: return + from hermes_constants import hermes_home_key + home_key = hermes_home_key() + # Timer threads start with empty ContextVars; without the scheduling turn's context the push would + # resolve the launch profile's home and credentials instead of the writing profile's. + ctx = _ctxvars.copy_context() def _fire(): with suppress(Exception): from tools.skills_sync_client import maybe_push_skills maybe_push_skills(message=f"sync: {skill_name}") with _sync_push_lock: - if _sync_push_timer is not None: - _sync_push_timer.cancel() # only sets an Event; never raises - _sync_push_timer = threading.Timer(_SYNC_PUSH_DEBOUNCE_S, _fire) - _sync_push_timer.daemon = True - _sync_push_timer.start() + pending = _sync_push_timers.get(home_key) + if pending is not None: + pending.cancel() # only sets an Event; never raises + timer = threading.Timer(_SYNC_PUSH_DEBOUNCE_S, ctx.run, args=(_fire,)) + timer.daemon = True + _sync_push_timers[home_key] = timer + timer.start() def _act_patch(a): diff --git a/tools/tirith_security.py b/tools/tirith_security.py index 79cee1db48..8814e17f2d 100644 --- a/tools/tirith_security.py +++ b/tools/tirith_security.py @@ -21,7 +21,7 @@ import time import urllib.request from contextlib import suppress -from hermes_constants import get_hermes_home +from hermes_constants import get_hermes_home, get_hermes_home_override, hermes_home_key logger = logging.getLogger(__name__) _REPO = "sheeki03/tirith" @@ -62,6 +62,9 @@ def _load_security_config() -> dict: _resolved_path: str | None | bool = None _INSTALL_FAILED = False _install_failure_reason: str = "" # reason tag when _resolved_path is _INSTALL_FAILED +# Routed profiles (multiplexed gateway) resolve their own binary: ``security.tirith_path`` and +# ``/bin/tirith`` are per profile, so the launch profile's slot above must not answer for them. +_resolved_path_by_home: dict[str, str] = {} # Circuit breaker: after _CRASH_LIMIT consecutive spawn/execution failures tirith is disabled # for the rest of the process so a broken binary can't turn every tool call into a fail-open @@ -108,12 +111,23 @@ def _warn_once(key: str, message: str, *args) -> None: def _cached_path() -> str | None: """The path resolved on a previous call, or None if unresolved (None) / failed (_INSTALL_FAILED).""" + if get_hermes_home_override() is not None: + return _resolved_path_by_home.get(hermes_home_key()) return _resolved_path or None +def _store_resolved(path: str) -> None: + global _resolved_path + if get_hermes_home_override() is not None: + _resolved_path_by_home[hermes_home_key()] = path + else: + _resolved_path = path + + def _set_resolved(path: str) -> None: - global _resolved_path, _install_failure_reason - _resolved_path, _install_failure_reason = path, "" + global _install_failure_reason + _store_resolved(path) + _install_failure_reason = "" def _set_failed(reason: str) -> None: @@ -361,7 +375,7 @@ def _resolve_locally(configured_path: str, *, warn_missing: bool) -> tuple[str | # An explicit (non-"tirith") path is authoritative: never auto-download a replacement. if configured_path != "tirith": if found := (expanded if _is_executable(expanded) else shutil.which(expanded)): - _resolved_path = found + _store_resolved(found) return found, False if warn_missing: logger.warning("Configured tirith path %r not found; scanning disabled", configured_path)