fix(multiplex): key tool-side and agent-side memos by profile home
Camofox VNC one-shot, computer-use aux-vision verdict, tirith binary path, MCP discovery lock path, remote-backend probe text, learned image token costs, auxiliary per-task semaphores and the custom-endpoint /models memo all held one profile's config-derived value for the whole process. The skill-sync debounce Timer ran with empty ContextVars, so a secondary's write pushed as the launch profile (and cancelled its pending push). Each memo is now keyed by hermes_home_key() (or credential fingerprint for the per-key catalog) under an override; the timer is per home and runs its callback inside the scheduling turn's copied context. Unscoped slots are unchanged.
This commit is contained in:
@@ -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.<task>.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:
|
||||
|
||||
+28
-14
@@ -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(
|
||||
|
||||
+20
-11
@@ -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]:
|
||||
|
||||
@@ -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 = ""
|
||||
|
||||
@@ -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"}
|
||||
@@ -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
|
||||
|
||||
@@ -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()))
|
||||
|
||||
+10
-4
@@ -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:
|
||||
|
||||
@@ -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):
|
||||
|
||||
@@ -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
|
||||
# ``<home>/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)
|
||||
|
||||
Reference in New Issue
Block a user