* feat(docker): publish nousresearch/hermes-sandbox:desktop for terminal backends
The terminal backends (docker, modal, daytona, singularity) all default to
nikolaik/python-nodejs:python3.11-nodejs20, a bare Python+Node base. For Bot
Screen, computer_use and the browser to run INSIDE that sandbox instead of on
the gateway host, the sandbox image needs the display stack.
docker/sandbox-desktop.Dockerfile is that base plus:
- the everyday tools it lacked (jq, ripgrep, fd, tmux, less, nano, vim,
zip, rsync, tree, procps, htop, sudo for the base's uid-1000 `pn`)
- the exact package set the Hermes -desktop image installs (TigerVNC,
Xfce components, dbus, xauth, fonts)
- Playwright's headed Chromium (same build as the -desktop image)
- cua-driver 0.28.2 from its pinned release tarball
No Hermes inside; the default user stays root like the base so nothing
changes for people who just switch docker_image. Desktop processes run as
`pn`. 4.27 GB on amd64.
docker.yml gains a `sandbox` variant with its own cache scope and repository
(nousresearch/hermes-sandbox:desktop, :main-desktop, :<release>-desktop);
the docker-integration suite is skipped for it (no Hermes to test) and
docker/sandbox-desktop-smoke.sh runs instead: as `pn`, every launcher and
cua-driver binary resolves, the real launcher.sh publishes :20, the RFB
socket completes the 3.8 handshake relayed over `docker exec -i` stdio, and
a headed Chromium maps a window on that display. hadolint lints the new
Dockerfile in docker-lint.yml.
* feat(docker): sandbox desktop base on python3.13-nodejs26
Matches the Hermes image (Python 3.13 / Node 26) and the top of requires-python;
the default docker_image tag it inherited was Python 3.11 / Node 20. Same pn
uid 1000, Debian 13; smoke (launcher, RFB relay, headed Chromium) passes.
* feat(docker): bake agent-browser into hermes-sandbox:desktop
The browser tools drive the agent-browser CLI; when the browser follows the
terminal backend that CLI has to exist inside the sandbox. Pinned to the same
^0.26.0 range the gateway resolves, --ignore-scripts like the gateway's npx path.
* feat(bot_desktop): the screen, computer_use and the browser follow the terminal backend
A user who sandboxes `terminal` (docker/ssh/singularity) had the agent's
screen, cua-driver and Chromium running on the gateway HOST beside that
sandbox: Bot Screen gave a headless host a display, the Xfce panel carries
xfce4-terminal, and `computer_use` could open a shell outside the boundary
the sandbox exists for.
Now the desktop lives where the terminal lives:
- tools/environments/streams.py: one primitive per spawn-per-call backend, the
local argv prefix that runs its remainder inside the sandbox with stdio open
(`docker exec -i`, `ssh`, `apptainer exec`). SDK backends (modal, daytona,
vercel) have none and report so.
- tools/bot_desktop/sandbox_host.py: launcher.sh runs inside the sandbox as
the image's `pn`; the pane's RFB bytes ride a 12-line python relay over that
prefix; `cua-driver mcp` is the prefix + the sandbox image's own driver.
- tools/bot_desktop/placement.py + `bot_desktop.placement` (auto|terminal|
gateway). `auto` follows the backend; a sandbox that cannot host a screen
REFUSES with the opt-in named instead of silently using the host.
- runtime.start/stop/status/published_env branch on placement; the pane,
lease, epoch fencing and CLI are unchanged.
- cua_backend: the MCP invocation is the sandbox one when the screen is
there; the host driver's runtime contract is irrelevant then; check_fn is
true under a terminal placement without a host binary.
- browser_tool_session: agent-browser invocations are wrapped in the prefix
with the daemon, socket dir and profile inside the sandbox; screenshots are
fetched back so MEDIA: paths keep working; recycle closes the sandbox
daemon.
- web_routers/display.py: the bridge pumps a relay's stdio when the screen is
in a sandbox, a unix socket otherwise.
Live on docker with nousresearch/hermes-sandbox:desktop: start/observe/RFB
handshake through the dashboard bridge, human takeover fences the agent
(HumanHasControl) and keystrokes reach the sandbox Xvnc, handback restores,
three start/stop rounds leave zero desktop processes; computer_use capture
and list_windows see only the sandbox's Xfce; browser_navigate/snapshot/
vision run with Chromium and agent-browser inside the container and zero
host processes on the bot profile; modal + auto refuses naming the opt-in.
* feat(desktop): Screen pane shows where a sandbox-placed screen runs; Install is host-only
DesktopStatus gains placement ('gateway' | 'terminal:<backend>'). A sandbox
image lacking the stack is a blocker naming hermes-sandbox:desktop, shown in
place of Start; install_command stays None there because the pane's Install
button runs the package manager on the gateway host, the wrong machine, and
display.install refuses for the same reason. The pane header carries
'Screen runs inside the docker sandbox, with the terminal' (4 locales).
* fix(bot_desktop): "is the screen in the sandbox" is a disk check on hot paths, never a config read
Every browser command and CUA spawn asked in_sandbox(), which resolves placement by
loading config, which initializes HERMES_HOME. Under a test's fake home that raised
HomeInitializationError from _run_browser_command; on a real host it read config per
click. Hot paths now ask sandbox_screen_running(): the start marker on disk, written
only by a sandbox start. Policy (in_sandbox) stays for start/install, where config is
the question. display.observe gates on "an RFB endpoint exists" for either placement.
* test(moa): late-accounting sink test asserts the wedged slot's row, not sink order
Under CI load the poll loop can see the interrupt before collecting the fast slot, so the
fast slot also arrives late and first; the test then failed on sink_calls[0]. The
contract is that the wedged slot's real usage reaches the sink.
* chore: retrigger CI (zero-job dispatch failure, auto-heal)
* fix(bot_desktop): a sandbox that died under a live screen fails loudly, never falls to the host
sandbox_screen_running() drops a start marker whose terminal environment is no longer
registered (stale after a process restart). When the environment object outlives its
container, the browser's sandbox wrap now checks the published DISPLAY and raises
"the screen inside the terminal backend's sandbox is gone; start it again" instead of
KeyError('AGENT_BROWSER_PROFILE'). Live: fresh sandbox navigate ok; docker rm -f the
container; next navigate returns that error; zero host Chromium either way.
* feat(terminal): nousresearch/hermes-sandbox:desktop is the default container sandbox
Every container backend (docker, modal, daytona, singularity) now defaults to the
sandbox image with the desktop stack, so Bot Screen, computer_use and the browser
run inside the sandbox for everyone who never chose an image; Python 3.13 / Node 26
match the Hermes image. One constant (DEFAULT_SANDBOX_IMAGE) replaces six copies of
the old literal. Migration 47 moves saved configs still holding the OLD default and
never touches an image the user pinned. Docker reuse recreates a container built
from another image, or the flip would silently never take effect for anyone with a
persisted container (live: old container removed, new one on 3.13 / Node 26 with
Xvnc, cua-driver, agent-browser present).
* chore: retrigger CI (zero-job dispatch failure)
* chore: retrigger CI (zero-job dispatch failure, auto-heal)
* chore: retrigger CI (zero-job dispatch failure, auto-heal)
* chore(config): template stamps v47 and shows the new default sandbox image
The template is what install.sh / docker / doctor --fix seed; a stamp behind
DEFAULT_CONFIG makes every fresh install migrate on first run.
* feat(sandbox): the default image change is a decision, not a surprise
A persisted Docker sandbox on another image is kept when docker_image is unset;
only a written docker_image (an explicit pin) recreates it. The pin verdict
travels as TERMINAL_DOCKER_IMAGE_PINNED through both terminal bridges (process
env and per-profile scope) and the container-config allowlist.
Approval surfaces, all through hermes_cli.sandbox_image_switch: the interactive
CLI asks once at startup (y = pin the new image, n = pin the current one, Enter =
ask later); the Screen pane shows the same choice with Switch / Keep buttons via
display.switchSandboxImage; `hermes config set terminal.docker_image …` is the
same answer from any shell. Gateways and cron never decide: they keep the sandbox
and log the notice.
Migration 47 now unsets a saved image equal to the OLD default instead of
rewriting it to the new one — that value was the template copied, not a pin, and
rewriting it would have made the runtime recreate existing sandboxes unasked.
Modal restores its snapshot and Daytona reuses its labeled sandbox regardless of
the configured image, so existing sandboxes there were already untouched.
* fix(config): both plain defaults that preceded the desktop sandbox image are template copies
main pinned nikolaik/python-nodejs:python3.14-nodejs22 (cd0f97f833) without a migration;
a saved config holding either literal is unset by migration 47, so it follows the default
and existing sandboxes get the keep-or-switch decision instead of a silent recreate.
* ci(docker): build the sandbox image on release/dispatch, not every main push
Leaves docker.yml exactly as on main. The sandbox image carries no Hermes code,
so two 4 GB multi-arch builds per merge bought nothing. sandbox-image.yml builds
and smokes on a PR that edits its own Dockerfile/smoke, and publishes only on a
release or a manual dispatch with publish=true. Stable tag stays :desktop.
* fix(config): keep main's config.py/config_defaults.py edits under the sandbox-image delta
The rebase resolved both files wholesale with the branch side, dropping main's move to
hermes_yaml (the 3.14 runtime venv has no PyYAML) and the 3.14 base pin. This is main's
version plus exactly the branch's own changes: DEFAULT_SANDBOX_IMAGE, the pin verdict in the
env bridge, placement defaults and the v47 stamp.
* test: sandbox-image tests read config.yaml through hermes_yaml (no PyYAML on the 3.14 runtime)
* fix(bot_desktop): read the sandbox marker BOM-tolerantly (windows footgun lint)
* fix(bot_desktop): placement is the authority; sandbox screen survives restarts
Review findings on the sandbox-hosted Bot Screen, each reproduced live first.
Authority. The browser preflight and the CUA invocation keyed off screen
LIVENESS, so `placement: terminal` with the screen not yet up handed the tool an
unchanged host command. `runtime.tool_placement()` is now the one resolver:
terminal placement starts the sandbox screen on demand (no auto_start opt-in
inside the user's own sandbox), refused placement raises its reason, and neither
ever yields the host. placement.resolve() answers a local backend from env alone
so the common case costs no config load on the spawn path.
Restart. sandbox_screen_running() deleted the marker whenever the process-local
terminal registry was empty, i.e. after every gateway restart, while Xvnc kept
running in the container; stop() then returned False and left it. The marker
now records the owning container; liveness comes from `docker inspect` on it,
stop/status re-attach to the recorded owner (even after the placement setting
moved), and only a container that is gone drops the marker.
SSH. remote_argv emitted `bash -c <script>` as three words; OpenSSH joins them
and the remote login shell ran `bash -c export` and the rest itself. The script
travels as one quoted word for ssh (remote_command knows the backend); docker
and apptainer keep argv.
CDP reach. agent-browser inside the sandbox reports the sandbox's loopback;
the Browser Use harness, browser_exec and the vault supervisor connect from the
host and got connection refused. streams.forward_port() proxies a local port
over the exec stream (same relay as the RFB bridge) and the CDP URL is rewritten
to the local end.
pids limit. --pids-limit 256 counts threads; measured on the desktop image the
desktop stack is 44, one Chromium tab 212, the agent's browser with two tabs
488. Past the cap every further docker exec died with "procReady not received".
Default is 2048 with the measurements in the comment.
Replacement. An approved image switch force-removed the old container before
`docker run` tried the new image; a private tag or registry outage left nothing.
The image is inspected/pulled first and the old container kept on failure.
Desktop integration. The sandbox start never passed the dock's browser launcher
(no Browser icon) and the thumbnail needed a host launcher pid + host ImageGrab
(always None). The dock runs the sandbox's Playwright Chromium on the shared
profile; the thumbnail is grabbed inside the sandbox (Pillow baked into the
image). The browser profile moves from /tmp — a 512 MB tmpfs emptied on every
container stop — to the desktop user's home, so logins follow the container.
Pin provenance. A TERMINAL_DOCKER_IMAGE written in a routed profile's .env is a
pin even when it spells the default; the scope compared values before.
* docs(bot-screen): no literal tmp path in the profile-location note
* fix(bot_desktop): docker inspect liveness probe closes stdin (TUI subprocess guard)
* fix(bot_desktop): adopting a screen the sandbox kept records the marker
Live ssh probe: after the host's state was lost while the sandbox kept its
Xvnc, start() took the idempotent early return (display already published)
and never wrote the host marker, so status/thumbnail/stop lost the screen.
Record the adopted display like a fresh launch.
Docs: what an ssh host of your own must carry, and why a Dockerfile ENV is
not enough for a login session (PLAYWRIGHT_BROWSERS_PATH via /etc/environment).
* docker(sandbox-desktop): login sessions find the browser (PLAYWRIGHT_BROWSERS_PATH via /etc/environment)
* docs(bot-screen): what the Apptainer path inherits from the image and what it does not
* chore(config): sandbox-image migration is 47→48 (main took 47 for compression.threshold_tokens)
* chore: retrigger CI (zero-job dispatch failure, auto-heal)
522 lines
29 KiB
Python
522 lines
29 KiB
Python
"""cua-driver MCP session plumbing: the asyncio bridge thread and the lazily-started, self-healing
|
|
``_CuaDriverSession`` (MCP transport with a ``cua-driver call`` CLI fallback). Config/policy
|
|
helpers are looked up lazily through the facade."""
|
|
|
|
from __future__ import annotations
|
|
|
|
import asyncio
|
|
import base64
|
|
import concurrent.futures
|
|
import contextlib
|
|
import json
|
|
import logging
|
|
import os
|
|
import threading
|
|
from typing import Any, Dict, List, Optional
|
|
|
|
from hermes_cli._subprocess_compat import windows_hide_flags
|
|
from tools.computer_use import cua_backend_driver as _driver
|
|
from tools.computer_use.cua_backend_parse import _extract_tool_result, _mcp_field, _tool_envelope
|
|
|
|
logger = logging.getLogger("tools.computer_use.cua_backend")
|
|
|
|
|
|
class _AsyncBridge:
|
|
"""Runs one asyncio loop on a daemon thread; marshals coroutines from the caller."""
|
|
|
|
def __init__(self) -> None:
|
|
self._loop: Optional[asyncio.AbstractEventLoop] = None
|
|
self._thread: Optional[threading.Thread] = None
|
|
self._ready = threading.Event()
|
|
|
|
def start(self) -> None:
|
|
if self._thread and self._thread.is_alive():
|
|
return
|
|
self._ready.clear()
|
|
|
|
def _run() -> None:
|
|
self._loop = asyncio.new_event_loop()
|
|
asyncio.set_event_loop(self._loop)
|
|
self._ready.set()
|
|
try:
|
|
self._loop.run_forever()
|
|
finally:
|
|
with contextlib.suppress(Exception):
|
|
self._loop.close()
|
|
|
|
self._thread = threading.Thread(target=_run, daemon=True, name="cua-driver-loop")
|
|
self._thread.start()
|
|
if not self._ready.wait(timeout=5.0):
|
|
raise RuntimeError("cua-driver asyncio bridge failed to start")
|
|
|
|
def run(self, coro, timeout: Optional[float] = 30.0) -> Any:
|
|
from agent.async_utils import safe_schedule_threadsafe
|
|
alive = self._loop is not None and self._thread is not None and self._thread.is_alive()
|
|
fut = safe_schedule_threadsafe(coro, self._loop) if alive else None # closes the coroutine on failure
|
|
if fut is None:
|
|
if asyncio.iscoroutine(coro):
|
|
coro.close() # no-op when safe_schedule_threadsafe already closed it
|
|
raise RuntimeError("cua-driver bridge not started")
|
|
return fut.result(timeout=timeout)
|
|
|
|
def stop(self) -> None:
|
|
if self._loop and self._loop.is_running():
|
|
self._loop.call_soon_threadsafe(self._loop.stop)
|
|
if self._thread:
|
|
self._thread.join(timeout=2.0)
|
|
self._thread = self._loop = None
|
|
|
|
# Fail-closed messages for calls whose effect on the remote screen is unknown. The action MAY have landed, so it
|
|
# is never replayed; the caller decides after taking fresh state.
|
|
_UNKNOWN_OUTCOME_MESSAGES = {
|
|
"transport_outcome_unknown": (
|
|
"cua-driver transport failed during {name}; the action outcome is unknown, so Hermes "
|
|
"did not replay it. Take fresh state before deciding whether to act again."),
|
|
"timeout_outcome_unknown": (
|
|
"cua-driver MCP call {name} timed out; the action outcome is unknown and may still have "
|
|
"taken effect on the remote screen. The session has been marked suspect and will be "
|
|
"recreated before the next computer-use call. Take fresh state before deciding "
|
|
"whether to act again."),
|
|
}
|
|
|
|
def _outcome_unknown(name: str, exc: Exception, code: str) -> Dict[str, Any]:
|
|
"""Fail-closed ``isError`` result for *code* (see ``_UNKNOWN_OUTCOME_MESSAGES``)."""
|
|
message = _UNKNOWN_OUTCOME_MESSAGES[code].format(name=name)
|
|
return _tool_envelope(message, [], {"ok": False, "code": code, "message": message, "operation": name,
|
|
"next_step": "fresh_state", "detail": str(exc)}, True, [])
|
|
|
|
def _tool_field(obj: Any, *names: str) -> Any:
|
|
"""``_mcp_field`` plus the ``model_extra`` fallback some MCP SDKs (Pydantic v2) forward custom fields via."""
|
|
value = _mcp_field(obj, names[0], names[-1])
|
|
return (getattr(obj, "model_extra", None) or {}).get(names[-1]) if value is None else value
|
|
|
|
_CLI_ATTEMPTS = 4 # CLI fallback transport retries (backoff 0.5s doubling)
|
|
|
|
def _cli_run_json(cmd: List[str], env: Dict[str, str], name: str, timeout: float) -> Any:
|
|
"""Run ``cua-driver call`` with backoff until it prints JSON; return the parsed value. "daemon is not running"
|
|
is PERMANENT for this invocation (the CLI needs the machine-wide daemon socket, which Linux installs typically
|
|
never start) -> fail fast, no ~3.5s backoff."""
|
|
import subprocess as _subprocess
|
|
import time as _time
|
|
|
|
backoff, last_err = 0.5, ""
|
|
for attempt in range(_CLI_ATTEMPTS):
|
|
try:
|
|
proc = _subprocess.run(cmd, capture_output=True, text=True, encoding="utf-8", errors="replace",
|
|
timeout=max(15.0, timeout), creationflags=windows_hide_flags(), env=env,
|
|
stdin=_subprocess.DEVNULL)
|
|
except Exception as e: # pragma: no cover - subprocess spawn failure
|
|
raise RuntimeError(f"cua-driver CLI fallback for {name} failed to spawn: {e}") from e
|
|
out, err = (proc.stdout or "").strip(), proc.stderr or ""
|
|
last_err = out[:200] or err[:200]
|
|
if "daemon is not running" in out or "daemon is not running" in err:
|
|
raise RuntimeError(f"cua-driver CLI fallback for {name} unavailable: the "
|
|
"machine-wide cua-driver daemon is not running (the "
|
|
"CLI transport requires it; the MCP runtime does not).")
|
|
start = min((i for i in (out.find("{"), out.find("[")) if i != -1), default=-1)
|
|
with contextlib.suppress(json.JSONDecodeError):
|
|
if start != -1:
|
|
return json.loads(out[start:])
|
|
if attempt < _CLI_ATTEMPTS - 1: # no JSON (EAGAIN warning / empty) — retry with backoff
|
|
logger.warning("cua-driver CLI fallback for %s got no JSON (attempt %d/%d); "
|
|
"retrying in %.1fs", name, attempt + 1, _CLI_ATTEMPTS, backoff)
|
|
_time.sleep(backoff)
|
|
backoff *= 2
|
|
raise RuntimeError(f"cua-driver CLI fallback for {name} returned no JSON after "
|
|
f"{_CLI_ATTEMPTS} attempts: {last_err}")
|
|
|
|
def _cli_result(parsed: Any, shot_file: Optional[str]) -> Dict[str, Any]:
|
|
"""Remap a ``cua-driver call`` JSON body into the ``_extract_tool_result`` shape (no ``image_mime_types`` key)."""
|
|
if not isinstance(parsed, dict):
|
|
return _tool_envelope(None, [], None, False)
|
|
# In-band logical failures with exit 0 must still fail closed.
|
|
is_error = parsed.get("isError") is True or parsed.get("is_error") is True
|
|
shot = parsed.get("screenshot_png_b64")
|
|
# Otherwise the screenshot was routed to a file (ours or the daemon's choice).
|
|
fpath = parsed.get("screenshot_file_path") or shot_file
|
|
if not shot and fpath and os.path.exists(fpath):
|
|
try:
|
|
with open(fpath, "rb") as fh:
|
|
shot = base64.b64encode(fh.read()).decode("ascii")
|
|
except Exception as e:
|
|
logger.debug("cua-driver CLI fallback: failed reading %s: %s", fpath, e)
|
|
data: Any = parsed.get("tree_markdown")
|
|
if data is not None and parsed.get("element_count") is not None:
|
|
data = f"{parsed['element_count']} elements\n{data}"
|
|
return _tool_envelope(data, [shot] if shot else [], parsed, is_error)
|
|
|
|
|
|
def _logical_error_text(result: Dict[str, Any]) -> str:
|
|
"""Flatten a logical MCP error into text for narrow classification."""
|
|
chunks: List[str] = []
|
|
for value in (result.get("data"), result.get("structuredContent")):
|
|
if value is None:
|
|
continue
|
|
try:
|
|
chunks.append(value if isinstance(value, str) else json.dumps(value, sort_keys=True))
|
|
except (TypeError, ValueError):
|
|
chunks.append(str(value))
|
|
return "\n".join(chunks)
|
|
|
|
def _is_ended_session_result(result: Any) -> bool:
|
|
"""Recognise cua-driver's explicit recoverable ended-session result."""
|
|
if not isinstance(result, dict) or result.get("isError") is not True:
|
|
return False
|
|
message = _logical_error_text(result).lower()
|
|
return ("session" in message and "start_session" in message
|
|
and ("has ended" in message or "session ended" in message))
|
|
|
|
|
|
class _CuaDriverSession:
|
|
"""Holds the mcp ClientSession. Spawned lazily; re-entered on drop. Lifecycle ownership: one long-running
|
|
coroutine (`_lifecycle_coro`) opens the stdio_client + ClientSession contexts, populates capabilities, sets
|
|
`_ready_event`, waits on `_shutdown_event`, then closes the contexts — enter and exit in the SAME task, as
|
|
anyio's cancel-scope invariant requires (each `bridge.run(coro)` is a NEW task). Tool calls run in short-lived
|
|
tasks touching only the session object."""
|
|
|
|
# Handshake calls issued BY start()/stop() — exempt from call_tool's auto-restart guard, or start() would recurse.
|
|
_LIFECYCLE_CALLS = frozenset({"start_session", "end_session"})
|
|
# Idempotent reads, safe to replay after a broken transport. Mutations stay out: a lost response does not
|
|
# prove they failed.
|
|
_TRANSPORT_REPLAY_SAFE_TOOLS = frozenset({"get_cursor_position", "get_displays", "get_screen_size",
|
|
"get_window_state", "list_apps", "list_windows"})
|
|
# A timed-out MCP session is wedged for later calls, so it is recreated before the next non-lifecycle
|
|
# call_tool. Class-level default: tests that bypass __init__ see healthy.
|
|
# See #74799.
|
|
_timeout_suspect = False
|
|
|
|
def __init__(self, bridge: _AsyncBridge, embedded_daemon: Optional[Any] = None) -> None:
|
|
self._bridge, self._embedded_daemon, self._session = bridge, embedded_daemon, None
|
|
self._lock, self._started = threading.Lock(), False
|
|
# Per-tool capability-token sets from `tools/list` (read via supports_capability). Raw input schemas are
|
|
# the source of truth for action properties: 0.9-era drivers advertise delivery_mode in inputSchema
|
|
# without the ``input.delivery_mode`` token.
|
|
# Keys are tool names (e.g. "click", "get_window_state"); values are sets of capability strings
|
|
# (e.g. "accessibility.element_tokens", "input.keyboard.type.terminal_safe"). Empty until the
|
|
# session starts; consumers should call `supports_capability` rather than reading directly. See
|
|
# #47072.
|
|
self._capabilities: Dict[str, set] = {}
|
|
self._tool_schemas: Dict[str, Dict[str, Any]] = {}
|
|
self._capability_version, self._ready_event = "", threading.Event()
|
|
self._shutdown_event: Optional[asyncio.Event] = None # created on bridge loop
|
|
self._lifecycle_future = None # concurrent.futures.Future
|
|
self._setup_error: Optional[BaseException] = None
|
|
# Declared via start_session; revives an ended-session rejection non-re-entrantly.
|
|
# Stable driver-side identity declared through start_session. Used to revive a logical ended-session
|
|
# rejection without recursive call_tool re-entry or backend-owned state (#71166).
|
|
self._declared_session_id: Optional[str] = None
|
|
self._transport_generation, self._transport_reset_callback = 0, None
|
|
|
|
async def _lifecycle_coro(self) -> None:
|
|
"""Owns the stdio MCP contexts: open, signal ready, block on shutdown, clean up — all in one task."""
|
|
import time as _time
|
|
from mcp import ClientSession, StdioServerParameters
|
|
from mcp.client.stdio import stdio_client
|
|
from tools.computer_use import cua_backend as _cb
|
|
from tools.environments.local import _sanitize_subprocess_env
|
|
|
|
self._shutdown_event = asyncio.Event() # built on the loop's own thread
|
|
_t0 = _time.monotonic()
|
|
# Phase marker: the ready-timeout error reports HOW FAR a wedged startup got.
|
|
# Phase marker surfaced by the ready-timeout error (issue #57025): when startup wedges, the caller
|
|
# reports HOW FAR it got instead of an opaque "never reached ready".
|
|
self._startup_phase = "binary-check"
|
|
try:
|
|
driver_cmd = _driver.resolve_cua_driver_cmd()
|
|
if not driver_cmd and _cb.sandbox_mcp_invocation() is None:
|
|
raise RuntimeError(_driver.cua_driver_install_hint())
|
|
self._startup_phase = "manifest-discovery"
|
|
daemon = self._embedded_daemon
|
|
(command, args), child_env = (
|
|
(daemon.proxy_invocation(), daemon.child_env()) if daemon is not None
|
|
else _cb.sandbox_mcp_invocation() or (_driver._resolve_mcp_invocation(driver_cmd), _cb.cua_driver_child_env()))
|
|
_t_manifest = _time.monotonic()
|
|
# Telemetry policy first (default: disabled), then strip Hermes secrets.
|
|
params = StdioServerParameters(command=command, args=args, env=_sanitize_subprocess_env(child_env))
|
|
async with stdio_client(params) as (read, write):
|
|
self._startup_phase = "mcp-initialize"
|
|
async with ClientSession(read, write) as session:
|
|
await session.initialize()
|
|
_t_init = _time.monotonic()
|
|
# Capabilities BEFORE exposing the session: the first call sees them.
|
|
self._startup_phase = "capability-discovery"
|
|
await self._populate_capabilities(session)
|
|
self._session, self._startup_phase = session, "ready"
|
|
self._ready_event.set()
|
|
logger.info("cua-driver session ready in %.1fs (manifest=%.1fs, mcp_init=%.1fs)",
|
|
_time.monotonic() - _t0, _t_manifest - _t0, _t_init - _t_manifest)
|
|
await self._shutdown_event.wait()
|
|
except BaseException as e:
|
|
# Ordinary errors and anyio CancelledError alike: start() surfaces this.
|
|
self._setup_error = e
|
|
self._ready_event.set()
|
|
raise
|
|
finally:
|
|
# A session that dies for ANY reason must be re-enterable: the next call sees _started False and
|
|
# rebuilds. Atomic bool write — stop() may hold _lock.
|
|
self._session, self._started = None, False
|
|
|
|
# Reset _started so a session that dies for ANY reason (MCP connection drop, driver crash, unexpected
|
|
# coro exit) is re-enterable: the next start()/call sees _started False and rebuilds the session instead
|
|
# of hanging forever on a dead one via _require_started(). On the normal stop() path this is a harmless
|
|
# idempotent no-op (stop() already set it False). A plain bool write is atomic in CPython, so this is
|
|
# safe from the bridge-loop thread without taking self._lock (which stop() may hold while awaiting this
|
|
# coro's future). See #55048 Bug 1.
|
|
async def _populate_capabilities(self, session: Any) -> None:
|
|
"""Cache per-tool capability sets, input schemas and capability_version from tools/list. Soft
|
|
prerequisite: on failure the map stays empty (capability False)."""
|
|
self._capabilities, self._tool_schemas, self._capability_version = {}, {}, ""
|
|
try:
|
|
tools_list = await session.list_tools()
|
|
for tool in getattr(tools_list, "tools", []) or []:
|
|
tool_name = getattr(tool, "name", None)
|
|
if not isinstance(tool_name, str):
|
|
continue
|
|
caps, schema = _tool_field(tool, "capabilities"), _tool_field(tool, "input_schema", "inputSchema")
|
|
self._capabilities[tool_name] = (
|
|
{c for c in caps if isinstance(c, str)} if isinstance(caps, list) else set())
|
|
self._tool_schemas[tool_name] = dict(schema) if isinstance(schema, dict) else {}
|
|
# capability_version is a sibling of `tools` in tools/list (NOT in initialize).
|
|
cv = _tool_field(tools_list, "capability_version")
|
|
if isinstance(cv, str):
|
|
self._capability_version = cv
|
|
except Exception as e:
|
|
logger.debug("cua-driver tools/list capability discovery failed: %s", e)
|
|
|
|
def start(self) -> None:
|
|
with self._lock:
|
|
if not self._started:
|
|
self._bridge.start()
|
|
self._start_lifecycle_locked()
|
|
self._started = True
|
|
|
|
def _start_lifecycle_locked(self) -> None:
|
|
"""Spawn the lifecycle owner and wait for ready. Caller holds self._lock."""
|
|
self._ready_event = threading.Event()
|
|
self._setup_error = self._shutdown_event = None
|
|
# The future tracks the WHOLE lifecycle; readiness is signalled via _ready_event.
|
|
loop = self._bridge._loop
|
|
if loop is None:
|
|
raise RuntimeError("cua-driver bridge not started")
|
|
self._lifecycle_future = asyncio.run_coroutine_threadsafe(self._lifecycle_coro(), loop)
|
|
if not self._ready_event.wait(timeout=30.0):
|
|
self._signal_shutdown_locked()
|
|
# Surface which startup phase wedged (issue #57025) — "doctor passes but the wrapper times out"
|
|
# reports are undiagnosable from a bare "never reached ready".
|
|
from hermes_constants import display_hermes_home
|
|
raise RuntimeError(
|
|
f"cua-driver session never reached ready (timeout 30s; stuck in phase: "
|
|
f"{getattr(self, '_startup_phase', 'unknown')}). Run `hermes computer-use doctor` and check "
|
|
f"{display_hermes_home()}/logs/agent.log for the phase timings.")
|
|
if self._setup_error is not None:
|
|
raise RuntimeError(f"cua-driver session setup failed: {self._setup_error}") from self._setup_error
|
|
self._transport_generation += 1
|
|
if self._transport_generation > 1:
|
|
self._notify_transport_reset()
|
|
|
|
def stop(self) -> None:
|
|
with self._lock:
|
|
if self._started:
|
|
self._started = False
|
|
self._stop_lifecycle_locked()
|
|
|
|
def set_transport_reset_callback(self, callback: Any) -> None:
|
|
"""Register a synchronous cache invalidation hook for transport swaps."""
|
|
self._transport_reset_callback = callback
|
|
|
|
def _notify_transport_reset(self) -> None:
|
|
try:
|
|
if (callback := getattr(self, "_transport_reset_callback", None)) is not None:
|
|
callback()
|
|
except Exception as exc:
|
|
logger.debug("cua-driver transport reset callback failed: %s", exc)
|
|
|
|
def _stop_lifecycle_locked(self) -> None:
|
|
self._signal_shutdown_locked()
|
|
fut, self._lifecycle_future = self._lifecycle_future, None
|
|
try:
|
|
if fut is not None:
|
|
fut.result(timeout=5.0)
|
|
except concurrent.futures.TimeoutError:
|
|
logger.warning("cua-driver session shutdown timed out (5s)")
|
|
except Exception as e:
|
|
logger.warning("cua-driver shutdown error: %s", e)
|
|
|
|
def _signal_shutdown_locked(self) -> None:
|
|
"""Set the asyncio shutdown event from the caller's thread."""
|
|
loop, event = self._bridge._loop, self._shutdown_event
|
|
if loop is not None and event is not None and loop.is_running():
|
|
with contextlib.suppress(RuntimeError): # loop closed — nothing to signal
|
|
loop.call_soon_threadsafe(event.set)
|
|
|
|
async def _call_tool_async(self, name: str, args: Dict[str, Any]) -> Dict[str, Any]:
|
|
return _extract_tool_result(await self._session.call_tool(name, args))
|
|
|
|
# ── Capability detection ─────────────────────────────────────────
|
|
# See #47072.
|
|
def supports_capability(self, capability: str, tool: Optional[str] = None) -> bool:
|
|
"""Driver advertises *capability* for *tool* (or ANY tool). False before start.
|
|
|
|
capability token (trycua/cua#1961 capability vocabulary).
|
|
"""
|
|
caps = [self._capabilities.get(tool, set())] if tool is not None else self._capabilities.values()
|
|
return any(capability in c for c in caps)
|
|
|
|
def _has_tool(self, name: str) -> bool:
|
|
"""``tools/list`` advertised *name*. Routes capture() (PNG capture moved into ``get_window_state``).
|
|
False before discovery — callers treat that as "unknown"."""
|
|
return name in self._capabilities
|
|
|
|
def supports_input_property(self, tool: str, property_name: str) -> bool:
|
|
"""Live tools/list schema accepts *property_name* (fails closed; no version guessing)."""
|
|
schema = getattr(self, "_tool_schemas", {}).get(tool, {})
|
|
properties = schema.get("properties") if isinstance(schema, dict) else None
|
|
return isinstance(properties, dict) and property_name in properties
|
|
|
|
@property
|
|
def capabilities_discovered(self) -> bool:
|
|
"""tools/list populated the map; when False ``_has_tool`` is untrustworthy."""
|
|
return bool(self._capabilities)
|
|
|
|
@property
|
|
def capability_version(self) -> str:
|
|
"""Driver-advertised capability vocabulary version ("" on old builds)."""
|
|
return self._capability_version
|
|
|
|
# ── Error classification (instance-patchable seams; result-shape checks live at module level) ──
|
|
@staticmethod
|
|
def _is_closed_session_error(exc: Exception) -> bool:
|
|
"""True for MCP/stdio failures that are recoverable by reconnecting."""
|
|
name, module = exc.__class__.__name__, getattr(exc.__class__, "__module__", "")
|
|
return (name in {"ClosedResourceError", "BrokenResourceError", "EndOfStream"}
|
|
or (module.startswith("anyio") and "Resource" in name)
|
|
or isinstance(exc, (BrokenPipeError, EOFError)))
|
|
|
|
@staticmethod
|
|
def _is_transient_daemon_error(exc: Exception) -> bool:
|
|
"""Daemon-proxy EAGAIN congestion: on macOS the ``cua-driver mcp`` bridge uses a non-blocking unix socket
|
|
and heavy ops (``get_window_state``) fail with ``os error 35`` when its buffer is full. A retry succeeds,
|
|
so back off / fall back instead of surfacing an empty 0x0 capture."""
|
|
msg = str(exc)
|
|
return any(needle in msg for needle in ("Resource temporarily unavailable", "os error 35",
|
|
"daemon transport error", "daemon proxy"))
|
|
|
|
# ── Recovery ─────────────────────────────────────────────────────
|
|
def _redeclare_session(self, timeout: float, failure_msg: str) -> bool:
|
|
"""start_session with the declared id; log *failure_msg* and return False on rejection."""
|
|
session_id = self._declared_session_id
|
|
result = self._bridge.run(self._call_tool_async("start_session", {"session": session_id}), timeout=timeout)
|
|
if result.get("isError") is True:
|
|
logger.warning(failure_msg, session_id, _logical_error_text(result))
|
|
return result.get("isError") is not True
|
|
|
|
def _recreate_session(self, name: str, timeout: float, log_msg: str, *, restart: bool = True,
|
|
clear_timeout_suspect: bool = False) -> None:
|
|
"""Log *log_msg* (``%s`` = *name*), then either start() a dead session or (``restart``) tear
|
|
down and rebuild the MCP lifecycle under ``_lock`` with capabilities repopulated from scratch;
|
|
finally re-attach the declared public label inside the replacement private lifecycle."""
|
|
logger.warning(log_msg, name)
|
|
if not restart:
|
|
self.start()
|
|
else:
|
|
with self._lock:
|
|
try:
|
|
if self._started:
|
|
self._stop_lifecycle_locked()
|
|
except Exception as e:
|
|
logger.debug("cua-driver session cleanup before reconnect failed: %s", e)
|
|
self._started = False
|
|
self._capabilities, self._tool_schemas, self._capability_version = {}, {}, ""
|
|
self._start_lifecycle_locked()
|
|
self._started = True
|
|
if clear_timeout_suspect:
|
|
self._timeout_suspect = False
|
|
if getattr(self, "_declared_session_id", None):
|
|
self._redeclare_session(timeout, "cua-driver public session label %s could not be restored: %s")
|
|
|
|
def _call_tool_via_cli(self, name: str, args: Dict[str, Any], timeout: float) -> Dict[str, Any]:
|
|
"""Fallback transport: ``cua-driver call <tool> <json>`` subprocess. The MCP stdio bridge can persistently
|
|
fail heavy calls (``get_window_state``) with EAGAIN while the plain CLI, on its own daemon socket, keeps
|
|
working. Output is remapped to the ``_extract_tool_result`` shape. ``get_window_state`` routes its
|
|
screenshot to a temp file (``screenshot_out_file``) so the daemon returns a tiny JSON body, not the
|
|
multi-megabyte base64 blob that congests the socket; ``_cli_result`` reads it back."""
|
|
import tempfile as _tempfile
|
|
from tools.computer_use import cua_backend as _cb
|
|
from tools.environments.local import _sanitize_subprocess_env
|
|
|
|
call_args, shot_file = dict(args), None
|
|
if name == "get_window_state" and "screenshot_out_file" not in call_args:
|
|
fd, shot_file = _tempfile.mkstemp(prefix="cua_shot_", suffix=".png")
|
|
os.close(fd)
|
|
call_args["screenshot_out_file"] = shot_file
|
|
driver_command = _driver.resolve_cua_driver_cmd()
|
|
if not driver_command:
|
|
raise RuntimeError(_driver.cua_driver_install_hint())
|
|
child_env, socket_args = _cb.cua_driver_child_env(), []
|
|
daemon = getattr(self, "_embedded_daemon", None)
|
|
if daemon is not None:
|
|
driver_command, child_env = daemon.proxy_invocation()[0], daemon.child_env()
|
|
socket_args = ["--socket", daemon.socket_path]
|
|
cmd = [driver_command, "call", name, json.dumps(call_args), *socket_args]
|
|
try:
|
|
return _cli_result(_cli_run_json(cmd, _sanitize_subprocess_env(child_env), name, timeout), shot_file)
|
|
finally:
|
|
if shot_file and os.path.exists(shot_file):
|
|
with contextlib.suppress(OSError):
|
|
os.remove(shot_file)
|
|
|
|
def call_tool(self, name: str, args: Dict[str, Any], timeout: float = 30.0) -> Dict[str, Any]:
|
|
if name not in self._LIFECYCLE_CALLS:
|
|
# A prior MCP timeout marks the session suspect (possibly wedged): recreate it so one timeout never
|
|
# poisons the run. Healthy sessions are never restarted here.
|
|
if self._timeout_suspect:
|
|
self._recreate_session(
|
|
name, timeout, "cua-driver session suspect after earlier MCP timeout; recreating before %s",
|
|
clear_timeout_suspect=True)
|
|
# A prior session may have died (MCP drop / driver crash) and reset _started.
|
|
if not self._started:
|
|
self._recreate_session(
|
|
name, timeout, "cua-driver session not active on %s; (re)starting before call", restart=False)
|
|
if not self._started:
|
|
raise RuntimeError("cua-driver session not started")
|
|
try:
|
|
result = self._bridge.run(self._call_tool_async(name, args), timeout=timeout)
|
|
except concurrent.futures.TimeoutError as e:
|
|
# Fail closed: the action may have landed, so never replay it.
|
|
# MCP deadline hit (#74799): the session is suspect and must be recreated before the next call.
|
|
# Fail closed — the action may have taken effect on the remote screen, so never replay it here;
|
|
# surface the uncertainty instead (#74799).
|
|
self._timeout_suspect = True
|
|
logger.warning("cua-driver MCP timed out on %s; marking session suspect "
|
|
"for recreation before the next call", name)
|
|
return _outcome_unknown(name, e, "timeout_outcome_unknown")
|
|
except Exception as e:
|
|
if self._is_transient_daemon_error(e):
|
|
if name not in self._TRANSPORT_REPLAY_SAFE_TOOLS:
|
|
self._notify_transport_reset()
|
|
return _outcome_unknown(name, e, "transport_outcome_unknown")
|
|
logger.warning("cua-driver MCP transport failed on %s (%s); "
|
|
"falling back to CLI transport", name, e)
|
|
return self._call_tool_via_cli(name, args, timeout)
|
|
if not self._is_closed_session_error(e):
|
|
raise
|
|
self._recreate_session(name, timeout, "cua-driver MCP session closed during %s; reconnecting once")
|
|
if name not in self._TRANSPORT_REPLAY_SAFE_TOOLS:
|
|
return _outcome_unknown(name, e, "transport_outcome_unknown")
|
|
result = self._bridge.run(self._call_tool_async(name, args), timeout=timeout)
|
|
# Remember only a SUCCESSFULLY declared identity: no stale recovery state.
|
|
declared_id, ok = args.get("session"), result.get("isError") is not True
|
|
if name == "start_session" and ok and isinstance(declared_id, str) and declared_id:
|
|
self._declared_session_id = declared_id
|
|
if _is_ended_session_result(result):
|
|
# Revive the stable session and replay the rejected call once; a 2nd rejection surfaces as-is.
|
|
# Never re-runs lifecycle calls -> an end_session result is final.
|
|
session_id = self._declared_session_id
|
|
if session_id and name not in self._LIFECYCLE_CALLS:
|
|
logger.warning("cua-driver session %s ended during %s; reviving and retrying once", session_id, name)
|
|
if self._redeclare_session(timeout, "cua-driver session %s could not be revived: %s"):
|
|
result = self._bridge.run(self._call_tool_async(name, args), timeout=timeout)
|
|
elif name == "end_session" and ok and declared_id == self._declared_session_id:
|
|
self._declared_session_id = None
|
|
return result
|