From fdaaa87ea4b1e56e71876fadfd3a923f1cac9ca0 Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Wed, 2 Sep 2026 16:21:07 -0700 Subject: [PATCH] refactor(browser): move session/daemon command execution to tools/browser_tool_session.py, CDP override + supervisor lifecycle to tools/browser_tool_cdp.py, vision helpers to tools/browser_tool_vision.py --- tools/browser_tool.py | 1165 ++------------------------------- tools/browser_tool_cdp.py | 195 ++++++ tools/browser_tool_session.py | 821 +++++++++++++++++++++++ tools/browser_tool_vision.py | 181 +++++ 4 files changed, 1239 insertions(+), 1123 deletions(-) create mode 100644 tools/browser_tool_cdp.py create mode 100644 tools/browser_tool_session.py create mode 100644 tools/browser_tool_vision.py diff --git a/tools/browser_tool.py b/tools/browser_tool.py index 1f584f1060..3e73edacad 100644 --- a/tools/browser_tool.py +++ b/tools/browser_tool.py @@ -22,13 +22,13 @@ import atexit import json import logging import os -import subprocess -import shutil +import subprocess # noqa: F401 (tests patch tools.browser_tool.subprocess.Popen) +import shutil # noqa: F401 (tests patch tools.browser_tool.shutil.which) import sys import tempfile import threading import time -from typing import Dict, Any, Optional, List, Tuple, Union +from typing import Dict, Any, Optional, Union from pathlib import Path from agent.redact import redact_cdp_url from hermes_constants import ( # noqa: F401 (test-patchable surface, read via origin by sibling modules) @@ -40,7 +40,7 @@ from hermes_constants import ( # noqa: F401 (test-patchable surface, read via ) from utils import env_int, is_truthy_value # noqa: F401 (read via origin by sibling modules) from hermes_cli.config import DEFAULT_CONFIG, cfg_get -from hermes_cli._subprocess_compat import windows_hide_flags +from hermes_cli._subprocess_compat import windows_hide_flags # noqa: F401 (test-patchable; read via origin) def __getattr__(name: str): @@ -317,279 +317,44 @@ def _get_open_command_timeout(*, first_open: bool = False) -> int: return max(base, floor) -def _needs_chromium_sandbox_bypass() -> bool: - """Return True when Chromium needs --no-sandbox to start reliably.""" - if hasattr(os, "geteuid") and os.geteuid() == 0: - return True - if _running_in_docker(): - return True - userns_restrict = "/proc/sys/kernel/apparmor_restrict_unprivileged_userns" - try: - with open(userns_restrict, encoding="utf-8") as f: - if f.read().strip() == "1": - return True - except OSError: - pass - return False - - -def _apply_chromium_sandbox_args(browser_env: Dict[str, str]) -> None: - """Add required Chromium sandbox flags without overriding user settings.""" - if ( - "AGENT_BROWSER_ARGS" not in browser_env - and "AGENT_BROWSER_CHROME_FLAGS" not in browser_env - and _needs_chromium_sandbox_bypass() - ): - logger.debug( - "browser: sandbox bypass needed (root/docker/AppArmor userns) — " - "injecting --no-sandbox" - ) - browser_env["AGENT_BROWSER_ARGS"] = "--no-sandbox,--disable-dev-shm-usage" - - -def _read_command_output_files(stdout_path: str, stderr_path: str) -> tuple[str, str]: - """Best-effort read of agent-browser stdout/stderr temp files.""" - stdout = stderr = "" - for path, slot in ((stdout_path, "stdout"), (stderr_path, "stderr")): - try: - with open(path, "r", encoding="utf-8") as f: - text = f.read().strip() - except OSError: - continue - if slot == "stdout": - stdout = text - else: - stderr = text - return stdout, stderr - - -def _unlink_command_output_files(*paths: str) -> None: - for path in paths: - try: - os.unlink(path) - except OSError: - pass - - -def _format_browser_timeout_error( - command: str, - timeout: int, - stdout: str, - stderr: str, -) -> str: - """Build an actionable timeout message from captured daemon output.""" - parts = [f"Command timed out after {timeout} seconds"] - detail = (stderr or stdout or "").strip() - if detail: - parts.append(detail[:1500]) - - combined = f"{stderr}\n{stdout}".lower() - hints: list[str] = [] - if "sandbox" in combined: - hints.append( - "Chromium sandbox launch failed. Set AGENT_BROWSER_ARGS=" - "'--no-sandbox,--disable-dev-shm-usage' in your environment, " - "or run: npx agent-browser install --with-deps" - ) - elif command == "open" and _is_local_mode(): - if _running_in_docker(): - hints.append( - "The browser daemon may still be starting or Chromium may be " - "missing. Pull the latest image: " - "docker pull ghcr.io/nousresearch/hermes-agent:latest" - ) - else: - hints.append( - "The browser daemon may still be starting, or Chromium may be " - "missing system libraries. Install/repair with: " - "npx agent-browser install --with-deps " - "(or: npx playwright install --with-deps chromium)" - ) - if hints: - parts.extend(hints) - return "\n".join(parts) - +from tools.browser_tool_session import ( # noqa: F401 (re-exported; tests patch tools.browser_tool.) + _needs_chromium_sandbox_bypass, + _apply_chromium_sandbox_args, + _read_command_output_files, + _unlink_command_output_files, + _format_browser_timeout_error, + _agent_browser_argv, + _prepare_session_socket_dir, + _agent_browser_command_env, + _popen_agent_browser, + _create_local_session, + _create_lightpanda_session, + _local_backend_process_dead, + _create_cdp_session, + _create_cloud_session_or_fallback, + _create_session_for_key, + _get_session_info, + _discard_timed_out_browser_session, + _read_browser_daemon_pid, + _browser_daemon_responsive, + _handle_browser_command_timeout, + _interpret_browser_command_output, + _run_browser_command, +) def _get_vision_model() -> Optional[str]: """Model for browser_vision (screenshot analysis — multimodal).""" return os.getenv("AUXILIARY_VISION_MODEL", "").strip() or None -def _resolve_cdp_override(cdp_url: str) -> str: - """Normalize a user-supplied CDP endpoint into a concrete websocket URL. - - Full ``ws://.../devtools/browser/...`` endpoints pass through; HTTP - discovery roots and bare ``ws://host:port`` are resolved via - ``/json/version`` → ``webSocketDebuggerUrl`` (falls back to the raw value - with a warning if discovery fails). - """ - raw = (cdp_url or "").strip() - if not raw: - return "" - - lowered = raw.lower() - if "/devtools/browser/" in lowered: - return raw - - discovery_url = raw - if lowered.startswith(("ws://", "wss://")): - if raw.count(":") == 2 and raw.rstrip("/").rsplit(":", 1)[-1].isdigit() and "/" not in raw.split(":", 2)[-1]: - discovery_url = ("http://" if lowered.startswith("ws://") else "https://") + raw.split("://", 1)[1] - else: - return raw - - if discovery_url.lower().endswith("/json/version"): - version_url = discovery_url - else: - version_url = discovery_url.rstrip("/") + "/json/version" - - try: - import requests # lazy — shared module object, test patches still apply - - response = requests.get(version_url, timeout=10) - response.raise_for_status() - payload = response.json() - except Exception as exc: - logger.warning( - "Failed to resolve CDP endpoint %s via %s: %s", - _sanitize_url_for_logs(raw), - _sanitize_url_for_logs(version_url), - _sanitize_url_for_logs(exc), - ) - return raw - - ws_url = str(payload.get("webSocketDebuggerUrl") or "").strip() - if ws_url: - logger.info( - "Resolved CDP endpoint %s -> %s", - _sanitize_url_for_logs(raw), - _sanitize_url_for_logs(ws_url), - ) - return ws_url - - logger.warning( - "CDP discovery at %s did not return webSocketDebuggerUrl; using raw endpoint", - _sanitize_url_for_logs(version_url), - ) - return raw - - -def _get_cdp_override_raw() -> str: - """Return the *configured* CDP override without any network I/O. - - Precedence: ``BROWSER_CDP_URL`` env (live ``/browser connect`` override), - then ``browser.cdp_url`` in config.yaml. Callers that only need to know - *whether* an override exists (check_fn gates, ``_is_local_mode`` / - ``_is_local_backend``, ``hermes doctor``) MUST use this, not - :func:`_get_cdp_override`: that one does a 10s HTTP discovery, and a stale - ``cdp_url`` pointing at a dead Chrome would stall every startup's schema - build with no error — no side effects during schema build. - """ - env_override = os.environ.get("BROWSER_CDP_URL", "").strip() - if env_override: - return env_override - return _browser_cfg( - "cdp_url", "", lambda v: str(v or "").strip(), "browser.cdp_url from config" - ) - - -def _get_cdp_override() -> str: - """Return the resolved CDP URL override, or "" (skips cloud AND local launch). - - May perform an HTTP ``/json/version`` discovery request — only call on - paths about to *connect* (session creation, supervisor attach); pure - is-it-configured gates must use :func:`_get_cdp_override_raw`. - """ - raw = _get_cdp_override_raw() - if not raw: - return "" - return _resolve_cdp_override(raw) - - -def _get_dialog_policy_config() -> Tuple[str, float]: - """Read ``browser.dialog_policy`` + ``browser.dialog_timeout_s`` from config. - - Returns a ``(policy, timeout_s)`` tuple, falling back to the supervisor's - defaults when keys are absent or invalid. - """ - # Defer imports so browser_tool can be imported in minimal environments. - from tools.browser_supervisor import ( - DEFAULT_DIALOG_POLICY, - DEFAULT_DIALOG_TIMEOUT_S, - _VALID_POLICIES, - ) - - try: - from hermes_cli.config import read_raw_config - - cfg = read_raw_config() - browser_cfg = cfg.get("browser", {}) if isinstance(cfg, dict) else {} - if not isinstance(browser_cfg, dict): - return DEFAULT_DIALOG_POLICY, DEFAULT_DIALOG_TIMEOUT_S - policy = str(browser_cfg.get("dialog_policy") or DEFAULT_DIALOG_POLICY) - if policy not in _VALID_POLICIES: - logger.debug("Invalid browser.dialog_policy=%r; using default", policy) - policy = DEFAULT_DIALOG_POLICY - timeout_raw = browser_cfg.get("dialog_timeout_s") - try: - timeout_s = float(timeout_raw) if timeout_raw is not None else DEFAULT_DIALOG_TIMEOUT_S - if timeout_s <= 0: - timeout_s = DEFAULT_DIALOG_TIMEOUT_S - except (TypeError, ValueError): - timeout_s = DEFAULT_DIALOG_TIMEOUT_S - return policy, timeout_s - except Exception: - return DEFAULT_DIALOG_POLICY, DEFAULT_DIALOG_TIMEOUT_S - - -def _ensure_cdp_supervisor(task_id: str) -> None: - """Start a CDP supervisor for ``task_id`` if an endpoint is reachable. - - Idempotent (``SupervisorRegistry.get_or_start`` skips an existing - ``(task_id, cdp_url)`` and restarts on URL change), so safe on every - navigate / ``/browser connect``. URL precedence: the CDP override, then the - session's own ``cdp_url`` (cloud providers). Swallows all errors — a failed - attach must not break the session; snapshots just lack - ``pending_dialogs`` / ``frame_tree``. - """ - cdp_url = _get_cdp_override() - if not cdp_url: - # Fallback: active session may carry a per-session CDP URL from a - # cloud provider (Browserbase sets this). - with _cleanup_lock: - session_info = _active_sessions.get(task_id, {}) - maybe = str(session_info.get("cdp_url") or "") - if maybe: - cdp_url = _resolve_cdp_override(maybe) - if not cdp_url: - return - try: - from tools.browser_supervisor import SUPERVISOR_REGISTRY # type: ignore[import-not-found] - - policy, timeout_s = _get_dialog_policy_config() - SUPERVISOR_REGISTRY.get_or_start( - task_id=task_id, - cdp_url=cdp_url, - dialog_policy=policy, - dialog_timeout_s=timeout_s, - ) - except Exception as exc: - logger.debug( - "CDP supervisor attach for task=%s failed (non-fatal): %s", - task_id, - exc, - ) - - -def _stop_cdp_supervisor(task_id: str) -> None: - """Stop the CDP supervisor for ``task_id`` if one exists. No-op otherwise.""" - try: - from tools.browser_supervisor import SUPERVISOR_REGISTRY # type: ignore[import-not-found] - - SUPERVISOR_REGISTRY.stop(task_id) - except Exception as exc: - logger.debug("CDP supervisor stop for task=%s failed (non-fatal): %s", task_id, exc) - +from tools.browser_tool_cdp import ( # noqa: F401 (re-exported; tests patch tools.browser_tool.) + _resolve_cdp_override, + _get_cdp_override_raw, + _get_cdp_override, + _get_dialog_policy_config, + _ensure_cdp_supervisor, + _stop_cdp_supervisor, +) # ============================================================================ # Cloud Provider Registry @@ -709,88 +474,6 @@ from tools.browser_tool_real_profile import ( # noqa: F401 ) -def _agent_browser_argv(browser_cmd: str) -> list: - """Command prefix to invoke agent-browser (concrete binary or npx sentinel). - - Concrete executable paths stay a single argv item (spaces intact); only the - synthetic npx sentinel expands. npx is resolved through the same - PATH + extended-PATH cascade ``_find_agent_browser`` uses — a bare - ``shutil.which("npx")`` would let a broken system npx shadow a healthy - Hermes-managed one. If npx isn't found at all (Termux, bare container) the - bare name is used so Popen raises a readable ``FileNotFoundError: 'npx'``. - ``--ignore-scripts``: AGENT_BROWSER_NPX_SPEC is a floating range, not an - exact pin — a compromised future patch must not run install-time scripts. - """ - if _is_npx_agent_browser_sentinel(browser_cmd): - _npx_bin = _resolve_npx_bin() or "npx" - return [_npx_bin, "--ignore-scripts", "--prefer-offline", "-y", AGENT_BROWSER_NPX_SPEC] - return [browser_cmd] - - -def _prepare_session_socket_dir(session_name: str) -> str: - """Create the per-session agent-browser socket dir and claim it with our PID. - - Each session gets its own dir so parallel workers don't fight over the - default socket path ("Failed to create socket directory: Permission - denied"). The owner_pid file is written BEFORE first use: another hermes - process's orphan reaper rmtree's any agent-browser-* dir in the shared - tmpdir that carries no live owner, which would delete this one mid-command. - """ - socket_dir = os.path.join(_socket_safe_tmpdir(), f"agent-browser-{session_name}") - os.makedirs(socket_dir, mode=0o700, exist_ok=True) - _write_owner_pid(socket_dir, session_name) - return socket_dir - - -def _agent_browser_command_env(socket_dir: str) -> Dict[str, str]: - """Credential-scrubbed env for one agent-browser command. - - Adds the discovery-time PATH fallbacks, the session socket dir, and the - daemon-side idle self-termination (``AGENT_BROWSER_IDLE_TIMEOUT_MS``, - agent-browser 0.24+) mirroring the Python-side inactivity janitor — - unless the user set the idle timeout explicitly. - """ - env = _build_browser_env() - env["PATH"] = _merge_browser_path(env.get("PATH", "")) - env["AGENT_BROWSER_SOCKET_DIR"] = socket_dir - if "AGENT_BROWSER_IDLE_TIMEOUT_MS" not in env: - env["AGENT_BROWSER_IDLE_TIMEOUT_MS"] = str(BROWSER_SESSION_INACTIVITY_TIMEOUT * 1000) - return env - - -def _popen_agent_browser(argv: List[str], env: Dict[str, str], socket_dir: str, tag: str) -> "subprocess.Popen": - """Spawn agent-browser with stdout/stderr redirected to ``socket_dir/_stdout_``. - - Temp files instead of pipes: the CLI forks a background daemon that inherits - its fds, so with pipes ``communicate()`` never sees EOF until the timeout. - Windows: CREATE_NO_WINDOW only (NOT CREATE_NEW_PROCESS_GROUP, which on - Python 3.11 cancels asyncio's running loop task and surfaces as - KeyboardInterrupt in the CLI), STARTF_USESTDHANDLES so CreateProcess hands - the child ONLY our three handles (leaked parent console handles make the - Rust binary's daemon grandchild die silently), close_fds=True for the rest. - Returns the Popen; the caller reads/unlinks the two files. - """ - stdout_path = os.path.join(socket_dir, f"_stdout_{tag}") - stderr_path = os.path.join(socket_dir, f"_stderr_{tag}") - stdout_fd = os.open(stdout_path, os.O_WRONLY | os.O_CREAT | os.O_TRUNC, 0o600) - stderr_fd = os.open(stderr_path, os.O_WRONLY | os.O_CREAT | os.O_TRUNC, 0o600) - try: - _popen_extra: dict = {} - if os.name == "nt": - _popen_extra["creationflags"] = windows_hide_flags() - _popen_extra["close_fds"] = True - _si = subprocess.STARTUPINFO() - _si.dwFlags |= subprocess.STARTF_USESTDHANDLES - _popen_extra["startupinfo"] = _si - return subprocess.Popen( - argv, stdout=stdout_fd, stderr=stderr_fd, - stdin=subprocess.DEVNULL, env=env, **_popen_extra, - ) - finally: - os.close(stdout_fd) - os.close(stderr_fd) - - def _url_is_private(url: str) -> bool: """Return True when the URL's host resolves to a private/LAN/loopback address. @@ -1268,250 +951,6 @@ BROWSER_TOOL_SCHEMAS = [ # Utility Functions # ============================================================================ -def _create_local_session(task_id: str, allow_real_profile: bool = True) -> Dict[str, str]: - import uuid - - # Real-profile consent: attach this local session (via CDP) to the user's - # browser running on a hermes-owned SNAPSHOT of their real profile, logins - # included. Fail closed on resolver/launch errors — a consented user must - # never be silently downgraded to a throwaway. The hybrid private-URL - # sidecar passes allow_real_profile=False: handing the user's cookie jar to - # an arbitrary internal host the model chose is a larger, unconsented - # exposure than the routing rule protects against (and a real-profile - # failure must not break private-URL routing). - if allow_real_profile: - cdp_url, err = _real_profile_cdp() - if err: - raise RuntimeError(err) - if cdp_url: - session_name = f"rp_{uuid.uuid4().hex[:10]}" - logger.info( - "Created real-profile local session %s for task %s", session_name, task_id - ) - return { - "session_name": session_name, - "bb_session_id": None, - "cdp_url": _resolve_cdp_override(cdp_url), - "features": {"local": True, "real_profile": True}, - } - - # Browser Use mode drives whatever CDP endpoint it is handed; with - # ``browser.engine: lightpanda`` that endpoint is a Hermes-spawned - # ``lightpanda serve``. The built-in tools never reach this branch — - # they are hidden in Browser Use mode — and keep driving Lightpanda via - # ``agent-browser --engine lightpanda`` on the plain local session below. - if _is_browser_use_cli_mode() and _using_lightpanda_engine(): - return _create_lightpanda_session(task_id) - - session_name = f"h_{uuid.uuid4().hex[:10]}" - logger.info("Created local browser session %s for task %s", - session_name, task_id) - return { - "session_name": session_name, - "bb_session_id": None, - "cdp_url": None, - "features": {"local": True}, - } - - -def _create_lightpanda_session(task_id: str) -> Dict[str, Any]: - """Spawn ``lightpanda serve`` for this session key (Browser Use mode).""" - import uuid - from tools.browser_lightpanda import launch_lightpanda - - session_name = f"lp_{uuid.uuid4().hex[:10]}" - server, err = launch_lightpanda( - session_name, block_private_networks=not _is_local_backend() - ) - if err: - raise RuntimeError(err) - logger.info( - "Created Lightpanda session %s (port %s) for task %s", - session_name, server.port, task_id, - ) - return { - "session_name": session_name, - "bb_session_id": None, - "cdp_url": server.cdp_url, - "features": {"local": True, "lightpanda": True}, - } - - -def _local_backend_process_dead(session_info: Dict[str, Any]) -> bool: - """True for a Lightpanda session whose ``lightpanda serve`` is gone.""" - if not (session_info.get("features") or {}).get("lightpanda"): - return False - from tools.browser_lightpanda import get_server - - server = get_server(session_info.get("session_name", "")) - return server is None or not server.is_alive() - - -def _create_cdp_session(task_id: str, cdp_url: str) -> Dict[str, str]: - """Create a session that connects to a user-supplied CDP endpoint.""" - import uuid - session_name = f"cdp_{uuid.uuid4().hex[:10]}" - logger.info("Created CDP browser session %s → %s for task %s", - session_name, _sanitize_url_for_logs(cdp_url), task_id) - return { - "session_name": session_name, - "bb_session_id": None, - "cdp_url": cdp_url, - "features": {"cdp_override": True}, - } - - -def _create_cloud_session_or_fallback(task_id: str, provider) -> Dict[str, Any]: - """Create a cloud session; fall back to local Chromium (marked degraded) on failure. - - Some cloud providers (Browser-Use v3) return an HTTP CDP discovery URL - instead of a raw websocket endpoint, so ``cdp_url`` is resolved here. - """ - try: - session_info = provider.create_session(task_id) - if not session_info or not isinstance(session_info, dict): - raise ValueError(f"Cloud provider returned invalid session: {session_info!r}") - if session_info.get("cdp_url"): - session_info = dict(session_info) - session_info["cdp_url"] = _resolve_cdp_override(str(session_info["cdp_url"])) - return session_info - except Exception as e: - provider_name = type(provider).__name__ - logger.warning( - "Cloud provider %s failed (%s); attempting fallback to local " - "Chromium for task %s", - provider_name, e, task_id, - exc_info=True, - ) - try: - session_info = _create_local_session(task_id) - except Exception as local_error: - raise RuntimeError( - f"Cloud provider {provider_name} failed ({e}) and local " - f"fallback also failed ({local_error})" - ) from e - # Mark session as degraded for observability - if isinstance(session_info, dict): - session_info = dict(session_info) - session_info["fallback_from_cloud"] = True - session_info["fallback_reason"] = str(e) - session_info["fallback_provider"] = provider_name - return session_info - - -def _create_session_for_key(task_id: str, force_local: bool) -> Dict[str, Any]: - """Create a fresh session for ``task_id`` (runs OUTSIDE the lock: cloud mode makes a network call). - - Precedence: CDP override > hybrid local sidecar > cloud provider > local. - The hybrid private-URL sidecar NEVER gets the real profile — presenting real - cookies to an arbitrary LAN host the model routed there is unconsented - exposure (see ``_create_local_session``). - """ - cdp_override = _get_cdp_override() - if cdp_override and not force_local: - return _create_cdp_session(task_id, cdp_override) - if force_local: - return _create_local_session(task_id, allow_real_profile=False) - provider = _get_cloud_provider() - if provider is None: - return _create_local_session(task_id) - return _create_cloud_session_or_fallback(task_id, provider) -def _get_session_info(task_id: Optional[str] = None) -> Dict[str, Any]: - """Get or create session info for a session key (thread-safe). - - ``task_id`` may carry the ``::local`` suffix (hybrid local sidecar), which - forces a local Chromium even when a cloud provider is configured. Also - starts the inactivity cleanup thread and touches activity tracking. - Returns a dict with ``session_name`` (always) plus ``bb_session_id`` / - ``cdp_url`` for cloud sessions. - """ - if task_id is None: - task_id = "default" - - # Start the cleanup thread if not running (handles inactivity timeouts) - _start_browser_cleanup_thread() - - # Update activity timestamp for this session - _update_session_activity(task_id) - - with _cleanup_lock: - # Check if we already have a session for this task - existing_session = _active_sessions.get(task_id) - - # Suspect-session recycle: a previous command - # timeout marked this cached session suspect via the SuspectableBackend - # adapter. ensure_healthy() tears it down here, at next use, and we fall - # through to create a fresh session — the expensive recycle lives on this - # path, not on the timeout path (mark must stay cheap). - if existing_session is not None and not _browser_session_backend(task_id).ensure_healthy(): - # Teardown removes the activity entry; the replacement must be - # tracked by the inactivity reaper like an initial session. - _update_session_activity(task_id) - with _cleanup_lock: - replacement = _active_sessions.get(task_id) - if replacement is not None and replacement is not existing_session: - # Another thread already recycled and re-created it. - return replacement - existing_session = None - - if existing_session is not None: - if ( - not _session_has_expired(existing_session) - and not _local_backend_process_dead(existing_session) - ): - return existing_session - - logger.info( - "Replacing expired or dead browser session for task %s", - task_id, - ) - _cleanup_single_browser_session(task_id) - # Cleanup removes the activity entry. The replacement session must be - # tracked by the inactivity reaper just like an initial session. - _update_session_activity(task_id) - - # Guard against a concurrent replacement: another thread may have - # already cleaned up the expired session and created a fresh one - # while we were waiting. If so, return the live replacement instead - # of falling through to create yet another session. - with _cleanup_lock: - replacement = _active_sessions.get(task_id) - if replacement is not None and replacement is not existing_session: - return replacement - - # Hybrid routing: session keys ending with ``::local`` force a local - # Chromium regardless of the globally-configured cloud provider. Public - # URLs in the same conversation continue to use the cloud session under - # the bare task_id key. - force_local = _is_local_sidecar_key(task_id) - session_info = _create_session_for_key(task_id, force_local) - - with _cleanup_lock: - # Double-check: another thread may have created a session while we - # were doing the network call. Use the existing one to avoid leaking - # orphan cloud sessions. - if task_id in _active_sessions: - return _active_sessions[task_id] - session_info = dict(session_info) - session_info.setdefault("session_key", task_id) - session_info.setdefault("owner_task_id", _bare_task_id_for_session_key(task_id)) - _active_sessions[task_id] = session_info - # A brand-new session is healthy by definition — drop any stale - # suspect flag left by a wedged-path eviction of its predecessor. - _suspect_browser_sessions.pop(task_id, None) - - # Lazy-start the CDP supervisor now that the session exists (if the - # backend surfaces a CDP URL via override or session_info["cdp_url"]). - # Idempotent; swallows errors. See _ensure_cdp_supervisor for details. - # Skip for local sidecars — they have no CDP URL — and for Lightpanda - # sessions: those only exist in Browser Use mode, where the browser_* - # tools that consume supervisor state are hidden, so the supervisor - # would just hold an idle second CDP connection to the process. - if not force_local and not (session_info.get("features") or {}).get("lightpanda"): - _ensure_cdp_supervisor(task_id) - - return session_info - from tools.browser_tool_snapshot import ( # noqa: F401 _store_full_snapshot, @@ -1521,370 +960,6 @@ from tools.browser_tool_snapshot import ( # noqa: F401 ) -def _discard_timed_out_browser_session( - task_id: str, - session_info: Dict[str, Any], - task_socket_dir: str, -) -> None: - """Drop a stuck client generation without losing cloud cleanup state.""" - with _cleanup_lock: - if _active_sessions.get(task_id) is not session_info: - return - _stop_cdp_supervisor(task_id) - if session_info.get("bb_session_id") or session_info.get("cdp_url"): - import uuid - replacement = dict(session_info) - replacement["session_name"] = f"h_{uuid.uuid4().hex[:10]}" - replacement.pop("_first_nav", None) - _active_sessions[task_id] = replacement - else: - _active_sessions.pop(task_id, None) - _session_last_activity.pop(task_id, None) - - bare_task_id = _bare_task_id_for_session_key(task_id) - if _last_active_session_key.get(bare_task_id) == task_id: - _last_active_session_key.pop(bare_task_id, None) - - session_name = str(session_info.get("session_name") or "") - if session_name: - pid_file = os.path.join(task_socket_dir, f"{session_name}.pid") - if os.path.isfile(pid_file): - try: - daemon_pid = int(Path(pid_file).read_text(encoding="utf-8").strip()) - if not _verify_reapable_browser_daemon(daemon_pid, task_socket_dir, session_name): - return - # Tree-kill: the daemon spawns Chromium - # children; terminating only the daemon PID leaks the whole - # Chromium tree. agent.deadline.kill_process_tree escalates - # SIGTERM → SIGKILL across the tree. - from agent import deadline as _deadline - - _deadline.kill_process_tree(daemon_pid) - except (ProcessLookupError, ValueError, PermissionError, OSError): - logger.debug("Could not kill timed-out browser daemon for %s", session_name) - return - shutil.rmtree(task_socket_dir, ignore_errors=True) - - -def _read_browser_daemon_pid(task_socket_dir: str, session_name: str) -> Optional[int]: - """Read the agent-browser daemon PID for a session (best-effort).""" - pid_file = os.path.join(task_socket_dir, f"{session_name}.pid") - try: - return int(Path(pid_file).read_text(encoding="utf-8").strip()) - except (OSError, ValueError): - return None - - -def _browser_daemon_responsive(task_socket_dir: str, probe_timeout_s: float = 1.0) -> bool: - """Cheap liveness probe: connect to the daemon's unix control socket. - - A successful connect proves the accept loop is alive (the command wedged on - the page/CDP side, not the daemon). Windows uses named pipes — no probe is - possible, so report unresponsive (tree-kill + respawn is the safe recovery). - """ - if os.name == "nt": - return False - import socket as socket_mod - - if not hasattr(socket_mod, "AF_UNIX"): - return False - try: - entries = os.listdir(task_socket_dir) - except OSError: - return False - sock_paths = [ - os.path.join(task_socket_dir, e) for e in entries if e.endswith(".sock") - ] - for sock_path in sock_paths: - try: - with socket_mod.socket(socket_mod.AF_UNIX, socket_mod.SOCK_STREAM) as s: - s.settimeout(probe_timeout_s) - s.connect(sock_path) - return True - except OSError: - continue - return False - - -def _handle_browser_command_timeout( - task_id: str, - session_info: Dict[str, Any], - task_socket_dir: str, -) -> None: - """Recover session state after a browser command timeout. - - * Cloud / CDP sessions: no local daemon to probe — replace the stuck client - generation now (fresh ``session_name``, same ``bb_session_id`` so cloud - cleanup still works). - * Local daemon alive (PID live, identity-verified, control socket accepts): - only the *command* wedged; mark the session suspect and let the next use - recycle it via ``ensure_healthy`` → clean ``close`` → fresh session. - * Local daemon wedged/dead: it cannot service a clean close and its Chromium - children would leak — tree-kill and evict now; the next call respawns. - - Both local branches ``mark_suspect`` first (cheap, lock-free) so the - poisoned-cache invariant holds even if eviction races another thread's - replacement (the flag then costs one harmless no-op teardown). - """ - if session_info.get("bb_session_id") or session_info.get("cdp_url"): - _discard_timed_out_browser_session(task_id, session_info, task_socket_dir) - return - - _browser_session_backend(task_id).mark_suspect( - "browser command timed out; session may be poisoned" - ) - - session_name = str(session_info.get("session_name") or "") - daemon_pid = _read_browser_daemon_pid(task_socket_dir, session_name) if session_name else None - daemon_alive = ( - daemon_pid is not None - and _pid_exists(daemon_pid) - and _verify_reapable_browser_daemon(daemon_pid, task_socket_dir, session_name) - and _browser_daemon_responsive(task_socket_dir) - ) - if daemon_alive: - logger.warning( - "browser daemon for %s is alive after command timeout; session " - "marked suspect and will be recycled at next use", task_id, - ) - return - - logger.warning( - "browser daemon for %s is wedged or dead after command timeout; " - "tree-killing and evicting the session", task_id, - ) - _discard_timed_out_browser_session(task_id, session_info, task_socket_dir) - # The poisoned entry is gone (evicted, or superseded by a concurrent - # replacement discard refused to touch) — either way the cache no longer - # holds the timed-out session, so drop the flag: it must not poison a - # session created later under the same key. - _suspect_browser_sessions.pop(task_id, None) - - -def _interpret_browser_command_output(command: str, stdout: str, stderr: str, returncode: int) -> Dict[str, Any]: - """Turn a finished agent-browser process's output into a result dict. - - Empty stdout with rc=0 is a broken state (stale daemon) and is reported as - failure rather than a silent ``{"success": True, "data": {}}`` — except for - commands in ``_EMPTY_OK_COMMANDS``. Non-JSON output is an error, except - for ``screenshot`` where the saved path is recovered from the prose. - """ - if stderr and stderr.strip(): - level = logging.WARNING if returncode != 0 else logging.DEBUG - logger.log(level, "browser '%s' stderr: %s", command, stderr.strip()[:500]) - - stdout_text = stdout.strip() - if not stdout_text and returncode == 0 and command not in _EMPTY_OK_COMMANDS: - logger.warning("browser '%s' returned empty output (rc=0)", command) - return {"success": False, "error": f"Browser command '{command}' returned no output"} - if not stdout_text: - if returncode != 0: - error_msg = stderr.strip() if stderr else f"Command failed with code {returncode}" - logger.warning("browser '%s' failed (rc=%s): %s", command, returncode, error_msg[:300]) - return {"success": False, "error": error_msg} - return {"success": True, "data": {}} - - try: - parsed = json.loads(stdout_text) - except json.JSONDecodeError: - raw = stdout_text[:2000] - logger.warning("browser '%s' returned non-JSON output (rc=%s): %s", - command, returncode, raw[:500]) - if command == "screenshot": - stderr_text = (stderr or "").strip() - combined_text = "\n".join(part for part in [stdout_text, stderr_text] if part) - recovered_path = _extract_screenshot_path_from_text(combined_text) - if recovered_path and Path(recovered_path).exists(): - logger.info( - "browser 'screenshot' recovered file from non-JSON output: %s", - recovered_path, - ) - return {"success": True, "data": {"path": recovered_path, "raw": raw}} - return {"success": False, "error": f"Non-JSON output from agent-browser for '{command}': {raw}"} - - # Empty snapshot content is a common sign of daemon/CDP issues. - if command == "snapshot" and parsed.get("success"): - snap_data = parsed.get("data", {}) - if not snap_data.get("snapshot") and not snap_data.get("refs"): - logger.warning("snapshot returned empty content. " - "Possible stale daemon or CDP connection issue. " - "returncode=%s", returncode) - return parsed - - -def _run_browser_command( - task_id: str, - command: str, - args: List[str] = None, - timeout: Optional[int] = None, - _engine_override: Optional[str] = None, -) -> Dict[str, Any]: - """Run one agent-browser CLI command against the task's session; returns its parsed JSON. - - ``timeout=None`` reads ``browser.command_timeout`` (default 30s). - ``_engine_override`` forces an engine for this call only (the Lightpanda - fallback uses it to retry with Chrome without touching global state). - """ - if timeout is None: - timeout = _safe_command_timeout() - args = args or [] - - # Build the command - try: - browser_cmd = _find_agent_browser() - except FileNotFoundError as e: - logger.warning("agent-browser CLI not found: %s", e) - return {"success": False, "error": str(e)} - - if _requires_real_termux_browser_install(browser_cmd): - error = _termux_browser_install_error() - logger.warning("browser command blocked on Termux: %s", error) - return {"success": False, "error": error} - - # Local mode with no Chromium on disk: fail fast with an actionable - # message instead of hanging for _command_timeout seconds per call. - # Skip when engine=lightpanda — LP doesn't need Chromium for navigation. - if ( - _is_local_mode() - and not _chromium_installed() - and _get_browser_engine() != "lightpanda" - and not _maybe_autoinstall_chromium() - ): - if _running_in_docker(): - hint = ( - "Chromium browser is missing. You're running in Docker — pull " - "the latest image to get the bundled Chromium: " - "docker pull ghcr.io/nousresearch/hermes-agent:latest" - ) - else: - hint = ( - "Chromium browser is missing. Install it with: " - "npx agent-browser install --with-deps " - "(or: npx playwright install --with-deps chromium)" - ) - logger.warning("browser command blocked: %s", hint) - return {"success": False, "error": hint} - - from tools.interrupt import is_interrupted - if is_interrupted(): - return {"success": False, "error": "Interrupted"} - - # Get session info (creates Browserbase session with proxies if needed) - try: - session_info = _get_session_info(task_id) - except Exception as e: - logger.warning("Failed to create browser session for task=%s: %s", task_id, e) - return {"success": False, "error": f"Failed to create browser session: {str(e)}"} - # Cleanup stops the supervisor before closing the backend; keep it stopped. - if command != "close" and session_info.get("cdp_url"): - _ensure_cdp_supervisor(task_id) - - # Build the command with the appropriate backend flag. - # Cloud mode: --cdp connects to Browserbase. - # Local mode: --session launches a local headless Chromium. - # The rest of the command (--json, command, args) is identical. - if session_info.get("cdp_url"): - # Cloud mode — connect to remote Browserbase browser via CDP - # IMPORTANT: Do NOT use --session with --cdp. In agent-browser >=0.13, - # --session creates a local browser instance and silently ignores --cdp. - backend_args = ["--cdp", session_info["cdp_url"]] - else: - # Local mode — launch Chromium (headless by default, headed when configured) - backend_args = ["--session", session_info["session_name"]] - if _is_headed_mode(): - backend_args.append("--headed") - - # Lightpanda engine injection (local mode only, agent-browser v0.25.3+). - # Use the resolved session backend rather than global cloud-provider state: - # hybrid private-URL routing can create a local sidecar while a cloud - # provider remains configured for public URLs. - engine = _engine_override or _get_browser_engine() - if engine != "auto" and not _is_camofox_mode() and not session_info.get("cdp_url"): - backend_args += ["--engine", engine] - - cmd_parts = _agent_browser_argv(browser_cmd) + backend_args + ["--json", command] + args - - try: - task_socket_dir = _prepare_session_socket_dir(session_info["session_name"]) - logger.debug("browser cmd=%s task=%s socket_dir=%s (%d chars)", - command, task_id, task_socket_dir, len(task_socket_dir)) - browser_env = _agent_browser_command_env(task_socket_dir) - - # Chromium-only launch flags are rejected by Lightpanda. Strip both - # the current and legacy variables for Lightpanda commands; explicit - # Chrome commands and fallback use the shared Chromium policy. - if engine == "lightpanda": - _stripped_args = browser_env.pop("AGENT_BROWSER_ARGS", None) - _stripped_flags = browser_env.pop("AGENT_BROWSER_CHROME_FLAGS", None) - if _stripped_args is not None or _stripped_flags is not None: - logger.debug( - "browser: stripped Chromium-only AGENT_BROWSER_ARGS/" - "AGENT_BROWSER_CHROME_FLAGS for Lightpanda command %s " - "(agent-browser rejects them with --engine lightpanda)", - command, - ) - else: - _apply_chromium_sandbox_args(browser_env) - - stdout_path = os.path.join(task_socket_dir, f"_stdout_{command}") - stderr_path = os.path.join(task_socket_dir, f"_stderr_{command}") - proc = _popen_agent_browser(cmd_parts, browser_env, task_socket_dir, command) - - try: - proc.wait(timeout=timeout) - except subprocess.TimeoutExpired: - proc.kill() - proc.wait() - stdout, stderr = _read_command_output_files(stdout_path, stderr_path) - _unlink_command_output_files(stdout_path, stderr_path) - _handle_browser_command_timeout(task_id, session_info, task_socket_dir) - if stderr and stderr.strip(): - logger.warning( - "browser '%s' stderr after timeout: %s", - command, - stderr.strip()[:500], - ) - logger.warning("browser '%s' timed out after %ds (task=%s, socket_dir=%s)", - command, timeout, task_id, task_socket_dir) - result = { - "success": False, - "error": _format_browser_timeout_error(command, timeout, stdout, stderr), - } - # Fall through to fallback check below - else: - with open(stdout_path, "r", encoding="utf-8") as f: - stdout = f.read() - with open(stderr_path, "r", encoding="utf-8") as f: - stderr = f.read() - _unlink_command_output_files(stdout_path, stderr_path) - result = _interpret_browser_command_output(command, stdout, stderr, proc.returncode) - - except Exception as e: - logger.warning("browser '%s' exception: %s", command, e, exc_info=True) - result = {"success": False, "error": str(e)} - - # --- Lightpanda automatic Chrome fallback --- - # If engine is lightpanda and the result looks broken, retry with Chrome. - # This runs for ALL exit paths (timeout, empty, non-JSON, nonzero rc, parsed). - fallback_reason = _lightpanda_fallback_reason(engine, command, result) - if fallback_reason: - logger.info( - "Lightpanda fallback: retrying '%s' with Chrome (task=%s): %s", - command, - task_id, - fallback_reason, - ) - # For screenshots, use the dedicated Chrome fallback helper - # (spins up a separate Chrome session to the same URL). - if command == "screenshot": - fallback_result = _chrome_fallback_screenshot(task_id, args or [], timeout) - else: - fallback_result = _run_chrome_fallback_command(task_id, command, args, timeout) - return _annotate_lightpanda_fallback(fallback_result, fallback_reason) - - return result - - # ============================================================================ # Browser Tool Functions # ============================================================================ @@ -2709,168 +1784,12 @@ _LP_VISION_FALLBACK_REASON = ( ) -def _vision_mode_label() -> str: - _cp = _get_cloud_provider() - return "local" if _cp is None else f"cloud ({_cp.provider_name()})" - - -def _lightpanda_vision_preroute( - effective_task_id: str, annotate: bool, screenshot_path: Path, -) -> Tuple[bool, Optional[str], Path]: - """Capture the vision screenshot through the Chrome fallback when Lightpanda is the engine. - - Lightpanda has no graphical renderer, so the normal path would fail with a - CDP error or return a placeholder PNG. Returns ``(prerouted, fallback_warning, - screenshot_path)``; on fallback failure ``prerouted`` is False and the caller - takes the normal screenshot path (forcing Chrome) so ``_run_browser_command`` - still produces the standard fallback metadata/error. - """ - engine = _get_browser_engine() - if engine != "lightpanda" or not _should_inject_engine(engine): - return False, None, screenshot_path - logger.debug("browser_vision: pre-routing screenshot to Chrome (engine=lightpanda)") - screenshot_args = ["--annotate"] if annotate else [] - fb_result = _chrome_fallback_screenshot(effective_task_id, screenshot_args, _get_command_timeout()) - fb_result = _annotate_lightpanda_fallback(fb_result, _LP_VISION_FALLBACK_REASON) - if not fb_result.get("success"): - logger.warning("Lightpanda Chrome fallback vision screenshot failed: %s", fb_result.get("error")) - return False, None, screenshot_path - fb_path = fb_result.get("data", {}).get("path", "") - if fb_path and os.path.exists(fb_path): - import uuid as uuid_mod - from hermes_constants import get_hermes_dir - - screenshots_dir = get_hermes_dir("cache/screenshots", "browser_screenshots") - screenshots_dir.mkdir(parents=True, exist_ok=True) - persistent_path = screenshots_dir / f"browser_screenshot_{uuid_mod.uuid4().hex}.png" - shutil.copy2(fb_path, persistent_path) - screenshot_path = persistent_path - return True, fb_result.get("fallback_warning"), screenshot_path - - -def _native_vision_result( - screenshot_path: Path, question: str, annotate: bool, - result: Dict[str, Any], lp_fallback_warning: Optional[str], -) -> Dict[str, Any]: - """Multimodal tool-result envelope: the main model inspects the pixels itself. - - History-reuse cap: this embed is baked into the tool result and re-sent on - every later turn, exactly like vision_analyze's native path — apply the same - proactive resize so full-res screenshots can't enter immutable history - uncapped. The helper's stat/dimension quick-estimate skips the resize when - already under both caps; without Pillow it fails open to the raw bytes. - """ - from tools.vision_tools import ( - _EMBED_MAX_DIMENSION, - _EMBED_TARGET_BYTES, - _build_native_vision_tool_result, - _resize_image_for_vision, - ) - - data_url = _resize_image_for_vision( - screenshot_path, - mime_type="image/png", - max_base64_bytes=_EMBED_TARGET_BYTES, - max_dimension=_EMBED_MAX_DIMENSION, - force_jpeg=True, - ) - native_result = _build_native_vision_tool_result( - image_url=str(screenshot_path), - question=question, - image_data_url=data_url, - image_size_bytes=screenshot_path.stat().st_size, - ) - meta = native_result.setdefault("meta", {}) - meta["screenshot_path"] = str(screenshot_path) - if lp_fallback_warning: - meta["fallback_warning"] = lp_fallback_warning - if annotate and result.get("data", {}).get("annotations"): - meta["annotations"] = result["data"]["annotations"] - native_result["text_summary"] = ( - f"{native_result.get('text_summary', '')} " - f"Screenshot path: {screenshot_path}" - ).strip() - return native_result - - -def _analyze_screenshot_with_aux_llm(screenshot_path: Path, question: str) -> str: - """One-shot aux vision-LLM analysis (not baked into history), secret-redacted. - - Encodes at full resolution; on a size-related provider rejection the image - is downscaled once and retried. Timeout/temperature come from - ``auxiliary.vision.*`` — local vision models (llama.cpp, ollama) can take - well over 30s, so the default timeout is generous. - """ - import base64 - - vision_prompt = ( - f"You are analyzing a screenshot of a web browser.\n\n" - f"User's question: {question}\n\n" - f"Provide a detailed and helpful answer based on what you see in the screenshot. " - f"If there are interactive elements, describe them. If there are verification challenges " - f"or CAPTCHAs, describe what type they are and what action might be needed. " - f"Focus on answering the user's specific question." - ) - _screenshot_bytes = screenshot_path.read_bytes() - _screenshot_b64 = base64.b64encode(_screenshot_bytes).decode("ascii") - data_url = f"data:image/png;base64,{_screenshot_b64}" - vision_model = _get_vision_model() - logger.debug("browser_vision: analysing screenshot (%d bytes)", - len(_screenshot_bytes)) - - vision_timeout = 120.0 - vision_temperature = 0.1 - try: - from hermes_cli.config import load_config - _vision_cfg = cfg_get(load_config(), "auxiliary", "vision", default={}) - _vt = _vision_cfg.get("timeout") - if _vt is not None: - vision_timeout = float(_vt) - _vtemp = _vision_cfg.get("temperature") - if _vtemp is not None: - vision_temperature = float(_vtemp) - except Exception: - pass - - call_kwargs = { - "task": "vision", - "messages": [ - { - "role": "user", - "content": [ - {"type": "text", "text": vision_prompt}, - {"type": "image_url", "image_url": {"url": data_url}}, - ], - } - ], - "temperature": vision_temperature, - "timeout": vision_timeout, - } - if vision_model: - call_kwargs["model"] = vision_model - try: - response = _lazy_call_llm(**call_kwargs) - except Exception as _api_err: - from tools.vision_tools import ( - _is_image_size_error, _resize_image_for_vision, _RESIZE_TARGET_BYTES, - ) - if not (_is_image_size_error(_api_err) and len(data_url) > _RESIZE_TARGET_BYTES): - raise - logger.info( - "Vision API rejected screenshot (%.1f MB); " - "auto-resizing to ~%.0f MB and retrying...", - len(data_url) / (1024 * 1024), - _RESIZE_TARGET_BYTES / (1024 * 1024), - ) - data_url = _resize_image_for_vision(screenshot_path, mime_type="image/png") - call_kwargs["messages"][0]["content"][1]["image_url"]["url"] = data_url - response = _lazy_call_llm(**call_kwargs) - - analysis = (response.choices[0].message.content or "").strip() - # Redact secrets the vision LLM may have read from the screenshot. - from agent.redact import redact_sensitive_text - return redact_sensitive_text(analysis) - +from tools.browser_tool_vision import ( # noqa: F401 (re-exported; tests patch tools.browser_tool.) + _vision_mode_label, + _lightpanda_vision_preroute, + _native_vision_result, + _analyze_screenshot_with_aux_llm, +) def browser_vision(question: str, annotate: bool = False, task_id: Optional[str] = None) -> Union[str, Dict[str, Any]]: """Screenshot the current page for visual inspection (CAPTCHAs, images, layouts). diff --git a/tools/browser_tool_cdp.py b/tools/browser_tool_cdp.py new file mode 100644 index 0000000000..1b98569a05 --- /dev/null +++ b/tools/browser_tool_cdp.py @@ -0,0 +1,195 @@ +"""User-supplied CDP endpoint resolution (browser.cdp_url / real-profile), dialog-policy config and the per-task CDP supervisor lifecycle. + +Split out of ``tools/browser_tool.py``; every name is re-imported there so +``tools.browser_tool.`` keeps resolving (and monkeypatching). Origin +symbols and module state are read/written through ``_bt`` (the origin module, +resolved per call by :func:`tools.browser_tool_origin.origin_module`) so +``patch("tools.browser_tool.X")`` is honoured and no import cycle exists. +""" + +import os +from typing import Tuple + +from tools.browser_tool_origin import origin_module as _origin + + +def _resolve_cdp_override(cdp_url: str) -> str: + """Normalize a user-supplied CDP endpoint into a concrete websocket URL. + + Full ``ws://.../devtools/browser/...`` endpoints pass through; HTTP + discovery roots and bare ``ws://host:port`` are resolved via + ``/json/version`` → ``webSocketDebuggerUrl`` (falls back to the raw value + with a warning if discovery fails). + """ + _bt = _origin() + raw = (cdp_url or "").strip() + if not raw: + return "" + + lowered = raw.lower() + if "/devtools/browser/" in lowered: + return raw + + discovery_url = raw + if lowered.startswith(("ws://", "wss://")): + if raw.count(":") == 2 and raw.rstrip("/").rsplit(":", 1)[-1].isdigit() and "/" not in raw.split(":", 2)[-1]: + discovery_url = ("http://" if lowered.startswith("ws://") else "https://") + raw.split("://", 1)[1] + else: + return raw + + if discovery_url.lower().endswith("/json/version"): + version_url = discovery_url + else: + version_url = discovery_url.rstrip("/") + "/json/version" + + try: + import requests # lazy — shared module object, test patches still apply + + response = requests.get(version_url, timeout=10) + response.raise_for_status() + payload = response.json() + except Exception as exc: + _bt.logger.warning( + "Failed to resolve CDP endpoint %s via %s: %s", + _bt._sanitize_url_for_logs(raw), + _bt._sanitize_url_for_logs(version_url), + _bt._sanitize_url_for_logs(exc), + ) + return raw + + ws_url = str(payload.get("webSocketDebuggerUrl") or "").strip() + if ws_url: + _bt.logger.info( + "Resolved CDP endpoint %s -> %s", + _bt._sanitize_url_for_logs(raw), + _bt._sanitize_url_for_logs(ws_url), + ) + return ws_url + + _bt.logger.warning( + "CDP discovery at %s did not return webSocketDebuggerUrl; using raw endpoint", + _bt._sanitize_url_for_logs(version_url), + ) + return raw + + +def _get_cdp_override_raw() -> str: + """Return the *configured* CDP override without any network I/O. + + Precedence: ``BROWSER_CDP_URL`` env (live ``/browser connect`` override), + then ``browser.cdp_url`` in config.yaml. Callers that only need to know + *whether* an override exists (check_fn gates, ``_is_local_mode`` / + ``_is_local_backend``, ``hermes doctor``) MUST use this, not + :func:`_get_cdp_override`: that one does a 10s HTTP discovery, and a stale + ``cdp_url`` pointing at a dead Chrome would stall every startup's schema + build with no error — no side effects during schema build. + """ + _bt = _origin() + env_override = os.environ.get("BROWSER_CDP_URL", "").strip() + if env_override: + return env_override + return _bt._browser_cfg( + "cdp_url", "", lambda v: str(v or "").strip(), "browser.cdp_url from config" + ) + + +def _get_cdp_override() -> str: + """Return the resolved CDP URL override, or "" (skips cloud AND local launch). + + May perform an HTTP ``/json/version`` discovery request — only call on + paths about to *connect* (session creation, supervisor attach); pure + is-it-configured gates must use :func:`_get_cdp_override_raw`. + """ + _bt = _origin() + raw = _bt._get_cdp_override_raw() + if not raw: + return "" + return _bt._resolve_cdp_override(raw) + + +def _get_dialog_policy_config() -> Tuple[str, float]: + """Read ``browser.dialog_policy`` + ``browser.dialog_timeout_s`` from config. + + Returns a ``(policy, timeout_s)`` tuple, falling back to the supervisor's + defaults when keys are absent or invalid. + """ + # Defer imports so browser_tool can be imported in minimal environments. + _bt = _origin() + from tools.browser_supervisor import ( + DEFAULT_DIALOG_POLICY, + DEFAULT_DIALOG_TIMEOUT_S, + _VALID_POLICIES, + ) + + try: + from hermes_cli.config import read_raw_config + + cfg = read_raw_config() + browser_cfg = cfg.get("browser", {}) if isinstance(cfg, dict) else {} + if not isinstance(browser_cfg, dict): + return DEFAULT_DIALOG_POLICY, DEFAULT_DIALOG_TIMEOUT_S + policy = str(browser_cfg.get("dialog_policy") or DEFAULT_DIALOG_POLICY) + if policy not in _VALID_POLICIES: + _bt.logger.debug("Invalid browser.dialog_policy=%r; using default", policy) + policy = DEFAULT_DIALOG_POLICY + timeout_raw = browser_cfg.get("dialog_timeout_s") + try: + timeout_s = float(timeout_raw) if timeout_raw is not None else DEFAULT_DIALOG_TIMEOUT_S + if timeout_s <= 0: + timeout_s = DEFAULT_DIALOG_TIMEOUT_S + except (TypeError, ValueError): + timeout_s = DEFAULT_DIALOG_TIMEOUT_S + return policy, timeout_s + except Exception: + return DEFAULT_DIALOG_POLICY, DEFAULT_DIALOG_TIMEOUT_S + + +def _ensure_cdp_supervisor(task_id: str) -> None: + """Start a CDP supervisor for ``task_id`` if an endpoint is reachable. + + Idempotent (``SupervisorRegistry.get_or_start`` skips an existing + ``(task_id, cdp_url)`` and restarts on URL change), so safe on every + navigate / ``/browser connect``. URL precedence: the CDP override, then the + session's own ``cdp_url`` (cloud providers). Swallows all errors — a failed + attach must not break the session; snapshots just lack + ``pending_dialogs`` / ``frame_tree``. + """ + _bt = _origin() + cdp_url = _bt._get_cdp_override() + if not cdp_url: + # Fallback: active session may carry a per-session CDP URL from a + # cloud provider (Browserbase sets this). + with _bt._cleanup_lock: + session_info = _bt._active_sessions.get(task_id, {}) + maybe = str(session_info.get("cdp_url") or "") + if maybe: + cdp_url = _bt._resolve_cdp_override(maybe) + if not cdp_url: + return + try: + from tools.browser_supervisor import SUPERVISOR_REGISTRY # type: ignore[import-not-found] + + policy, timeout_s = _bt._get_dialog_policy_config() + SUPERVISOR_REGISTRY.get_or_start( + task_id=task_id, + cdp_url=cdp_url, + dialog_policy=policy, + dialog_timeout_s=timeout_s, + ) + except Exception as exc: + _bt.logger.debug( + "CDP supervisor attach for task=%s failed (non-fatal): %s", + task_id, + exc, + ) + + +def _stop_cdp_supervisor(task_id: str) -> None: + """Stop the CDP supervisor for ``task_id`` if one exists. No-op otherwise.""" + _bt = _origin() + try: + from tools.browser_supervisor import SUPERVISOR_REGISTRY # type: ignore[import-not-found] + + SUPERVISOR_REGISTRY.stop(task_id) + except Exception as exc: + _bt.logger.debug("CDP supervisor stop for task=%s failed (non-fatal): %s", task_id, exc) diff --git a/tools/browser_tool_session.py b/tools/browser_tool_session.py new file mode 100644 index 0000000000..d988d9214c --- /dev/null +++ b/tools/browser_tool_session.py @@ -0,0 +1,821 @@ +"""agent-browser session management: daemon spawn, per-backend session creation (local/lightpanda/cdp/cloud), cached session lookup, command execution with timeout handling and output interpretation. + +Split out of ``tools/browser_tool.py``; every name is re-imported there so +``tools.browser_tool.`` keeps resolving (and monkeypatching). Origin +symbols and module state are read/written through ``_bt`` (the origin module, +resolved per call by :func:`tools.browser_tool_origin.origin_module`) so +``patch("tools.browser_tool.X")`` is honoured and no import cycle exists. +""" + +import json +import logging +import os +import shutil +import subprocess +from pathlib import Path +from typing import Any, Dict, List, Optional + +from tools.browser_tool_origin import origin_module as _origin + + +def _needs_chromium_sandbox_bypass() -> bool: + """Return True when Chromium needs --no-sandbox to start reliably.""" + _bt = _origin() + if hasattr(os, "geteuid") and os.geteuid() == 0: + return True + if _bt._running_in_docker(): + return True + userns_restrict = "/proc/sys/kernel/apparmor_restrict_unprivileged_userns" + try: + with open(userns_restrict, encoding="utf-8") as f: + if f.read().strip() == "1": + return True + except OSError: + pass + return False + + +def _apply_chromium_sandbox_args(browser_env: Dict[str, str]) -> None: + """Add required Chromium sandbox flags without overriding user settings.""" + _bt = _origin() + if ( + "AGENT_BROWSER_ARGS" not in browser_env + and "AGENT_BROWSER_CHROME_FLAGS" not in browser_env + and _bt._needs_chromium_sandbox_bypass() + ): + _bt.logger.debug( + "browser: sandbox bypass needed (root/docker/AppArmor userns) — " + "injecting --no-sandbox" + ) + browser_env["AGENT_BROWSER_ARGS"] = "--no-sandbox,--disable-dev-shm-usage" + + +def _read_command_output_files(stdout_path: str, stderr_path: str) -> tuple[str, str]: + """Best-effort read of agent-browser stdout/stderr temp files.""" + stdout = stderr = "" + for path, slot in ((stdout_path, "stdout"), (stderr_path, "stderr")): + try: + with open(path, "r", encoding="utf-8") as f: + text = f.read().strip() + except OSError: + continue + if slot == "stdout": + stdout = text + else: + stderr = text + return stdout, stderr + + +def _unlink_command_output_files(*paths: str) -> None: + for path in paths: + try: + os.unlink(path) + except OSError: + pass + + +def _format_browser_timeout_error( + command: str, + timeout: int, + stdout: str, + stderr: str, +) -> str: + """Build an actionable timeout message from captured daemon output.""" + _bt = _origin() + parts = [f"Command timed out after {timeout} seconds"] + detail = (stderr or stdout or "").strip() + if detail: + parts.append(detail[:1500]) + + combined = f"{stderr}\n{stdout}".lower() + hints: list[str] = [] + if "sandbox" in combined: + hints.append( + "Chromium sandbox launch failed. Set AGENT_BROWSER_ARGS=" + "'--no-sandbox,--disable-dev-shm-usage' in your environment, " + "or run: npx agent-browser install --with-deps" + ) + elif command == "open" and _bt._is_local_mode(): + if _bt._running_in_docker(): + hints.append( + "The browser daemon may still be starting or Chromium may be " + "missing. Pull the latest image: " + "docker pull ghcr.io/nousresearch/hermes-agent:latest" + ) + else: + hints.append( + "The browser daemon may still be starting, or Chromium may be " + "missing system libraries. Install/repair with: " + "npx agent-browser install --with-deps " + "(or: npx playwright install --with-deps chromium)" + ) + if hints: + parts.extend(hints) + return "\n".join(parts) + + +def _agent_browser_argv(browser_cmd: str) -> list: + """Command prefix to invoke agent-browser (concrete binary or npx sentinel). + + Concrete executable paths stay a single argv item (spaces intact); only the + synthetic npx sentinel expands. npx is resolved through the same + PATH + extended-PATH cascade ``_find_agent_browser`` uses — a bare + ``shutil.which("npx")`` would let a broken system npx shadow a healthy + Hermes-managed one. If npx isn't found at all (Termux, bare container) the + bare name is used so Popen raises a readable ``FileNotFoundError: 'npx'``. + ``--ignore-scripts``: AGENT_BROWSER_NPX_SPEC is a floating range, not an + exact pin — a compromised future patch must not run install-time scripts. + """ + _bt = _origin() + if _bt._is_npx_agent_browser_sentinel(browser_cmd): + _npx_bin = _bt._resolve_npx_bin() or "npx" + return [_npx_bin, "--ignore-scripts", "--prefer-offline", "-y", _bt.AGENT_BROWSER_NPX_SPEC] + return [browser_cmd] + + +def _prepare_session_socket_dir(session_name: str) -> str: + """Create the per-session agent-browser socket dir and claim it with our PID. + + Each session gets its own dir so parallel workers don't fight over the + default socket path ("Failed to create socket directory: Permission + denied"). The owner_pid file is written BEFORE first use: another hermes + process's orphan reaper rmtree's any agent-browser-* dir in the shared + tmpdir that carries no live owner, which would delete this one mid-command. + """ + _bt = _origin() + socket_dir = os.path.join(_bt._socket_safe_tmpdir(), f"agent-browser-{session_name}") + os.makedirs(socket_dir, mode=0o700, exist_ok=True) + _bt._write_owner_pid(socket_dir, session_name) + return socket_dir + + +def _agent_browser_command_env(socket_dir: str) -> Dict[str, str]: + """Credential-scrubbed env for one agent-browser command. + + Adds the discovery-time PATH fallbacks, the session socket dir, and the + daemon-side idle self-termination (``AGENT_BROWSER_IDLE_TIMEOUT_MS``, + agent-browser 0.24+) mirroring the Python-side inactivity janitor — + unless the user set the idle timeout explicitly. + """ + _bt = _origin() + env = _bt._build_browser_env() + env["PATH"] = _bt._merge_browser_path(env.get("PATH", "")) + env["AGENT_BROWSER_SOCKET_DIR"] = socket_dir + if "AGENT_BROWSER_IDLE_TIMEOUT_MS" not in env: + env["AGENT_BROWSER_IDLE_TIMEOUT_MS"] = str(_bt.BROWSER_SESSION_INACTIVITY_TIMEOUT * 1000) + return env + + +def _popen_agent_browser(argv: List[str], env: Dict[str, str], socket_dir: str, tag: str) -> "subprocess.Popen": + """Spawn agent-browser with stdout/stderr redirected to ``socket_dir/_stdout_``. + + Temp files instead of pipes: the CLI forks a background daemon that inherits + its fds, so with pipes ``communicate()`` never sees EOF until the timeout. + Windows: CREATE_NO_WINDOW only (NOT CREATE_NEW_PROCESS_GROUP, which on + Python 3.11 cancels asyncio's running loop task and surfaces as + KeyboardInterrupt in the CLI), STARTF_USESTDHANDLES so CreateProcess hands + the child ONLY our three handles (leaked parent console handles make the + Rust binary's daemon grandchild die silently), close_fds=True for the rest. + Returns the Popen; the caller reads/unlinks the two files. + """ + _bt = _origin() + stdout_path = os.path.join(socket_dir, f"_stdout_{tag}") + stderr_path = os.path.join(socket_dir, f"_stderr_{tag}") + stdout_fd = os.open(stdout_path, os.O_WRONLY | os.O_CREAT | os.O_TRUNC, 0o600) + stderr_fd = os.open(stderr_path, os.O_WRONLY | os.O_CREAT | os.O_TRUNC, 0o600) + try: + _popen_extra: dict = {} + if os.name == "nt": + _popen_extra["creationflags"] = _bt.windows_hide_flags() + _popen_extra["close_fds"] = True + _si = subprocess.STARTUPINFO() + _si.dwFlags |= subprocess.STARTF_USESTDHANDLES + _popen_extra["startupinfo"] = _si + return subprocess.Popen( + argv, stdout=stdout_fd, stderr=stderr_fd, + stdin=subprocess.DEVNULL, env=env, **_popen_extra, + ) + finally: + os.close(stdout_fd) + os.close(stderr_fd) + + +def _create_local_session(task_id: str, allow_real_profile: bool = True) -> Dict[str, str]: + _bt = _origin() + import uuid + + # Real-profile consent: attach this local session (via CDP) to the user's + # browser running on a hermes-owned SNAPSHOT of their real profile, logins + # included. Fail closed on resolver/launch errors — a consented user must + # never be silently downgraded to a throwaway. The hybrid private-URL + # sidecar passes allow_real_profile=False: handing the user's cookie jar to + # an arbitrary internal host the model chose is a larger, unconsented + # exposure than the routing rule protects against (and a real-profile + # failure must not break private-URL routing). + if allow_real_profile: + cdp_url, err = _bt._real_profile_cdp() + if err: + raise RuntimeError(err) + if cdp_url: + session_name = f"rp_{uuid.uuid4().hex[:10]}" + _bt.logger.info( + "Created real-profile local session %s for task %s", session_name, task_id + ) + return { + "session_name": session_name, + "bb_session_id": None, + "cdp_url": _bt._resolve_cdp_override(cdp_url), + "features": {"local": True, "real_profile": True}, + } + + # Browser Use mode drives whatever CDP endpoint it is handed; with + # ``browser.engine: lightpanda`` that endpoint is a Hermes-spawned + # ``lightpanda serve``. The built-in tools never reach this branch — + # they are hidden in Browser Use mode — and keep driving Lightpanda via + # ``agent-browser --engine lightpanda`` on the plain local session below. + if _bt._is_browser_use_cli_mode() and _bt._using_lightpanda_engine(): + return _bt._create_lightpanda_session(task_id) + + session_name = f"h_{uuid.uuid4().hex[:10]}" + _bt.logger.info("Created local browser session %s for task %s", + session_name, task_id) + return { + "session_name": session_name, + "bb_session_id": None, + "cdp_url": None, + "features": {"local": True}, + } + + +def _create_lightpanda_session(task_id: str) -> Dict[str, Any]: + """Spawn ``lightpanda serve`` for this session key (Browser Use mode).""" + _bt = _origin() + import uuid + from tools.browser_lightpanda import launch_lightpanda + + session_name = f"lp_{uuid.uuid4().hex[:10]}" + server, err = launch_lightpanda( + session_name, block_private_networks=not _bt._is_local_backend() + ) + if err: + raise RuntimeError(err) + _bt.logger.info( + "Created Lightpanda session %s (port %s) for task %s", + session_name, server.port, task_id, + ) + return { + "session_name": session_name, + "bb_session_id": None, + "cdp_url": server.cdp_url, + "features": {"local": True, "lightpanda": True}, + } + + +def _local_backend_process_dead(session_info: Dict[str, Any]) -> bool: + """True for a Lightpanda session whose ``lightpanda serve`` is gone.""" + if not (session_info.get("features") or {}).get("lightpanda"): + return False + from tools.browser_lightpanda import get_server + + server = get_server(session_info.get("session_name", "")) + return server is None or not server.is_alive() + + +def _create_cdp_session(task_id: str, cdp_url: str) -> Dict[str, str]: + """Create a session that connects to a user-supplied CDP endpoint.""" + _bt = _origin() + import uuid + session_name = f"cdp_{uuid.uuid4().hex[:10]}" + _bt.logger.info("Created CDP browser session %s → %s for task %s", + session_name, _bt._sanitize_url_for_logs(cdp_url), task_id) + return { + "session_name": session_name, + "bb_session_id": None, + "cdp_url": cdp_url, + "features": {"cdp_override": True}, + } + + +def _create_cloud_session_or_fallback(task_id: str, provider) -> Dict[str, Any]: + """Create a cloud session; fall back to local Chromium (marked degraded) on failure. + + Some cloud providers (Browser-Use v3) return an HTTP CDP discovery URL + instead of a raw websocket endpoint, so ``cdp_url`` is resolved here. + """ + _bt = _origin() + try: + session_info = provider.create_session(task_id) + if not session_info or not isinstance(session_info, dict): + raise ValueError(f"Cloud provider returned invalid session: {session_info!r}") + if session_info.get("cdp_url"): + session_info = dict(session_info) + session_info["cdp_url"] = _bt._resolve_cdp_override(str(session_info["cdp_url"])) + return session_info + except Exception as e: + provider_name = type(provider).__name__ + _bt.logger.warning( + "Cloud provider %s failed (%s); attempting fallback to local " + "Chromium for task %s", + provider_name, e, task_id, + exc_info=True, + ) + try: + session_info = _bt._create_local_session(task_id) + except Exception as local_error: + raise RuntimeError( + f"Cloud provider {provider_name} failed ({e}) and local " + f"fallback also failed ({local_error})" + ) from e + # Mark session as degraded for observability + if isinstance(session_info, dict): + session_info = dict(session_info) + session_info["fallback_from_cloud"] = True + session_info["fallback_reason"] = str(e) + session_info["fallback_provider"] = provider_name + return session_info + + +def _create_session_for_key(task_id: str, force_local: bool) -> Dict[str, Any]: + """Create a fresh session for ``task_id`` (runs OUTSIDE the lock: cloud mode makes a network call). + + Precedence: CDP override > hybrid local sidecar > cloud provider > local. + The hybrid private-URL sidecar NEVER gets the real profile — presenting real + cookies to an arbitrary LAN host the model routed there is unconsented + exposure (see ``_create_local_session``). + """ + _bt = _origin() + cdp_override = _bt._get_cdp_override() + if cdp_override and not force_local: + return _bt._create_cdp_session(task_id, cdp_override) + if force_local: + return _bt._create_local_session(task_id, allow_real_profile=False) + provider = _bt._get_cloud_provider() + if provider is None: + return _bt._create_local_session(task_id) + return _bt._create_cloud_session_or_fallback(task_id, provider) + + +def _get_session_info(task_id: Optional[str] = None) -> Dict[str, Any]: + """Get or create session info for a session key (thread-safe). + + ``task_id`` may carry the ``::local`` suffix (hybrid local sidecar), which + forces a local Chromium even when a cloud provider is configured. Also + starts the inactivity cleanup thread and touches activity tracking. + Returns a dict with ``session_name`` (always) plus ``bb_session_id`` / + ``cdp_url`` for cloud sessions. + """ + _bt = _origin() + if task_id is None: + task_id = "default" + + # Start the cleanup thread if not running (handles inactivity timeouts) + _bt._start_browser_cleanup_thread() + + # Update activity timestamp for this session + _bt._update_session_activity(task_id) + + with _bt._cleanup_lock: + # Check if we already have a session for this task + existing_session = _bt._active_sessions.get(task_id) + + # Suspect-session recycle: a previous command + # timeout marked this cached session suspect via the SuspectableBackend + # adapter. ensure_healthy() tears it down here, at next use, and we fall + # through to create a fresh session — the expensive recycle lives on this + # path, not on the timeout path (mark must stay cheap). + if existing_session is not None and not _bt._browser_session_backend(task_id).ensure_healthy(): + # Teardown removes the activity entry; the replacement must be + # tracked by the inactivity reaper like an initial session. + _bt._update_session_activity(task_id) + with _bt._cleanup_lock: + replacement = _bt._active_sessions.get(task_id) + if replacement is not None and replacement is not existing_session: + # Another thread already recycled and re-created it. + return replacement + existing_session = None + + if existing_session is not None: + if ( + not _bt._session_has_expired(existing_session) + and not _bt._local_backend_process_dead(existing_session) + ): + return existing_session + + _bt.logger.info( + "Replacing expired or dead browser session for task %s", + task_id, + ) + _bt._cleanup_single_browser_session(task_id) + # Cleanup removes the activity entry. The replacement session must be + # tracked by the inactivity reaper just like an initial session. + _bt._update_session_activity(task_id) + + # Guard against a concurrent replacement: another thread may have + # already cleaned up the expired session and created a fresh one + # while we were waiting. If so, return the live replacement instead + # of falling through to create yet another session. + with _bt._cleanup_lock: + replacement = _bt._active_sessions.get(task_id) + if replacement is not None and replacement is not existing_session: + return replacement + + # Hybrid routing: session keys ending with ``::local`` force a local + # Chromium regardless of the globally-configured cloud provider. Public + # URLs in the same conversation continue to use the cloud session under + # the bare task_id key. + force_local = _bt._is_local_sidecar_key(task_id) + session_info = _bt._create_session_for_key(task_id, force_local) + + with _bt._cleanup_lock: + # Double-check: another thread may have created a session while we + # were doing the network call. Use the existing one to avoid leaking + # orphan cloud sessions. + if task_id in _bt._active_sessions: + return _bt._active_sessions[task_id] + session_info = dict(session_info) + session_info.setdefault("session_key", task_id) + session_info.setdefault("owner_task_id", _bt._bare_task_id_for_session_key(task_id)) + _bt._active_sessions[task_id] = session_info + # A brand-new session is healthy by definition — drop any stale + # suspect flag left by a wedged-path eviction of its predecessor. + _bt._suspect_browser_sessions.pop(task_id, None) + + # Lazy-start the CDP supervisor now that the session exists (if the + # backend surfaces a CDP URL via override or session_info["cdp_url"]). + # Idempotent; swallows errors. See _ensure_cdp_supervisor for details. + # Skip for local sidecars — they have no CDP URL — and for Lightpanda + # sessions: those only exist in Browser Use mode, where the browser_* + # tools that consume supervisor state are hidden, so the supervisor + # would just hold an idle second CDP connection to the process. + if not force_local and not (session_info.get("features") or {}).get("lightpanda"): + _bt._ensure_cdp_supervisor(task_id) + + return session_info + + +def _discard_timed_out_browser_session( + task_id: str, + session_info: Dict[str, Any], + task_socket_dir: str, +) -> None: + """Drop a stuck client generation without losing cloud cleanup state.""" + _bt = _origin() + with _bt._cleanup_lock: + if _bt._active_sessions.get(task_id) is not session_info: + return + _bt._stop_cdp_supervisor(task_id) + if session_info.get("bb_session_id") or session_info.get("cdp_url"): + import uuid + replacement = dict(session_info) + replacement["session_name"] = f"h_{uuid.uuid4().hex[:10]}" + replacement.pop("_first_nav", None) + _bt._active_sessions[task_id] = replacement + else: + _bt._active_sessions.pop(task_id, None) + _bt._session_last_activity.pop(task_id, None) + + bare_task_id = _bt._bare_task_id_for_session_key(task_id) + if _bt._last_active_session_key.get(bare_task_id) == task_id: + _bt._last_active_session_key.pop(bare_task_id, None) + + session_name = str(session_info.get("session_name") or "") + if session_name: + pid_file = os.path.join(task_socket_dir, f"{session_name}.pid") + if os.path.isfile(pid_file): + try: + daemon_pid = int(Path(pid_file).read_text(encoding="utf-8").strip()) + if not _bt._verify_reapable_browser_daemon(daemon_pid, task_socket_dir, session_name): + return + # Tree-kill: the daemon spawns Chromium + # children; terminating only the daemon PID leaks the whole + # Chromium tree. agent.deadline.kill_process_tree escalates + # SIGTERM → SIGKILL across the tree. + from agent import deadline as _deadline + + _deadline.kill_process_tree(daemon_pid) + except (ProcessLookupError, ValueError, PermissionError, OSError): + _bt.logger.debug("Could not kill timed-out browser daemon for %s", session_name) + return + shutil.rmtree(task_socket_dir, ignore_errors=True) + + +def _read_browser_daemon_pid(task_socket_dir: str, session_name: str) -> Optional[int]: + """Read the agent-browser daemon PID for a session (best-effort).""" + pid_file = os.path.join(task_socket_dir, f"{session_name}.pid") + try: + return int(Path(pid_file).read_text(encoding="utf-8").strip()) + except (OSError, ValueError): + return None + + +def _browser_daemon_responsive(task_socket_dir: str, probe_timeout_s: float = 1.0) -> bool: + """Cheap liveness probe: connect to the daemon's unix control socket. + + A successful connect proves the accept loop is alive (the command wedged on + the page/CDP side, not the daemon). Windows uses named pipes — no probe is + possible, so report unresponsive (tree-kill + respawn is the safe recovery). + """ + if os.name == "nt": + return False + import socket as socket_mod + + if not hasattr(socket_mod, "AF_UNIX"): + return False + try: + entries = os.listdir(task_socket_dir) + except OSError: + return False + sock_paths = [ + os.path.join(task_socket_dir, e) for e in entries if e.endswith(".sock") + ] + for sock_path in sock_paths: + try: + with socket_mod.socket(socket_mod.AF_UNIX, socket_mod.SOCK_STREAM) as s: + s.settimeout(probe_timeout_s) + s.connect(sock_path) + return True + except OSError: + continue + return False + + +def _handle_browser_command_timeout( + task_id: str, + session_info: Dict[str, Any], + task_socket_dir: str, +) -> None: + """Recover session state after a browser command timeout. + + * Cloud / CDP sessions: no local daemon to probe — replace the stuck client + generation now (fresh ``session_name``, same ``bb_session_id`` so cloud + cleanup still works). + * Local daemon alive (PID live, identity-verified, control socket accepts): + only the *command* wedged; mark the session suspect and let the next use + recycle it via ``ensure_healthy`` → clean ``close`` → fresh session. + * Local daemon wedged/dead: it cannot service a clean close and its Chromium + children would leak — tree-kill and evict now; the next call respawns. + + Both local branches ``mark_suspect`` first (cheap, lock-free) so the + poisoned-cache invariant holds even if eviction races another thread's + replacement (the flag then costs one harmless no-op teardown). + """ + _bt = _origin() + if session_info.get("bb_session_id") or session_info.get("cdp_url"): + _bt._discard_timed_out_browser_session(task_id, session_info, task_socket_dir) + return + + _bt._browser_session_backend(task_id).mark_suspect( + "browser command timed out; session may be poisoned" + ) + + session_name = str(session_info.get("session_name") or "") + daemon_pid = _bt._read_browser_daemon_pid(task_socket_dir, session_name) if session_name else None + daemon_alive = ( + daemon_pid is not None + and _bt._pid_exists(daemon_pid) + and _bt._verify_reapable_browser_daemon(daemon_pid, task_socket_dir, session_name) + and _bt._browser_daemon_responsive(task_socket_dir) + ) + if daemon_alive: + _bt.logger.warning( + "browser daemon for %s is alive after command timeout; session " + "marked suspect and will be recycled at next use", task_id, + ) + return + + _bt.logger.warning( + "browser daemon for %s is wedged or dead after command timeout; " + "tree-killing and evicting the session", task_id, + ) + _bt._discard_timed_out_browser_session(task_id, session_info, task_socket_dir) + # The poisoned entry is gone (evicted, or superseded by a concurrent + # replacement discard refused to touch) — either way the cache no longer + # holds the timed-out session, so drop the flag: it must not poison a + # session created later under the same key. + _bt._suspect_browser_sessions.pop(task_id, None) + + +def _interpret_browser_command_output(command: str, stdout: str, stderr: str, returncode: int) -> Dict[str, Any]: + """Turn a finished agent-browser process's output into a result dict. + + Empty stdout with rc=0 is a broken state (stale daemon) and is reported as + failure rather than a silent ``{"success": True, "data": {}}`` — except for + commands in ``_EMPTY_OK_COMMANDS``. Non-JSON output is an error, except + for ``screenshot`` where the saved path is recovered from the prose. + """ + _bt = _origin() + if stderr and stderr.strip(): + level = logging.WARNING if returncode != 0 else logging.DEBUG + _bt.logger.log(level, "browser '%s' stderr: %s", command, stderr.strip()[:500]) + + stdout_text = stdout.strip() + if not stdout_text and returncode == 0 and command not in _bt._EMPTY_OK_COMMANDS: + _bt.logger.warning("browser '%s' returned empty output (rc=0)", command) + return {"success": False, "error": f"Browser command '{command}' returned no output"} + if not stdout_text: + if returncode != 0: + error_msg = stderr.strip() if stderr else f"Command failed with code {returncode}" + _bt.logger.warning("browser '%s' failed (rc=%s): %s", command, returncode, error_msg[:300]) + return {"success": False, "error": error_msg} + return {"success": True, "data": {}} + + try: + parsed = json.loads(stdout_text) + except json.JSONDecodeError: + raw = stdout_text[:2000] + _bt.logger.warning("browser '%s' returned non-JSON output (rc=%s): %s", + command, returncode, raw[:500]) + if command == "screenshot": + stderr_text = (stderr or "").strip() + combined_text = "\n".join(part for part in [stdout_text, stderr_text] if part) + recovered_path = _bt._extract_screenshot_path_from_text(combined_text) + if recovered_path and Path(recovered_path).exists(): + _bt.logger.info( + "browser 'screenshot' recovered file from non-JSON output: %s", + recovered_path, + ) + return {"success": True, "data": {"path": recovered_path, "raw": raw}} + return {"success": False, "error": f"Non-JSON output from agent-browser for '{command}': {raw}"} + + # Empty snapshot content is a common sign of daemon/CDP issues. + if command == "snapshot" and parsed.get("success"): + snap_data = parsed.get("data", {}) + if not snap_data.get("snapshot") and not snap_data.get("refs"): + _bt.logger.warning("snapshot returned empty content. " + "Possible stale daemon or CDP connection issue. " + "returncode=%s", returncode) + return parsed + + +def _run_browser_command( + task_id: str, + command: str, + args: List[str] = None, + timeout: Optional[int] = None, + _engine_override: Optional[str] = None, +) -> Dict[str, Any]: + """Run one agent-browser CLI command against the task's session; returns its parsed JSON. + + ``timeout=None`` reads ``browser.command_timeout`` (default 30s). + ``_engine_override`` forces an engine for this call only (the Lightpanda + fallback uses it to retry with Chrome without touching global state). + """ + _bt = _origin() + if timeout is None: + timeout = _bt._safe_command_timeout() + args = args or [] + + # Build the command + try: + browser_cmd = _bt._find_agent_browser() + except FileNotFoundError as e: + _bt.logger.warning("agent-browser CLI not found: %s", e) + return {"success": False, "error": str(e)} + + if _bt._requires_real_termux_browser_install(browser_cmd): + error = _bt._termux_browser_install_error() + _bt.logger.warning("browser command blocked on Termux: %s", error) + return {"success": False, "error": error} + + # Local mode with no Chromium on disk: fail fast with an actionable + # message instead of hanging for _command_timeout seconds per call. + # Skip when engine=lightpanda — LP doesn't need Chromium for navigation. + if ( + _bt._is_local_mode() + and not _bt._chromium_installed() + and _bt._get_browser_engine() != "lightpanda" + and not _bt._maybe_autoinstall_chromium() + ): + if _bt._running_in_docker(): + hint = ( + "Chromium browser is missing. You're running in Docker — pull " + "the latest image to get the bundled Chromium: " + "docker pull ghcr.io/nousresearch/hermes-agent:latest" + ) + else: + hint = ( + "Chromium browser is missing. Install it with: " + "npx agent-browser install --with-deps " + "(or: npx playwright install --with-deps chromium)" + ) + _bt.logger.warning("browser command blocked: %s", hint) + return {"success": False, "error": hint} + + from tools.interrupt import is_interrupted + if is_interrupted(): + return {"success": False, "error": "Interrupted"} + + # Get session info (creates Browserbase session with proxies if needed) + try: + session_info = _bt._get_session_info(task_id) + except Exception as e: + _bt.logger.warning("Failed to create browser session for task=%s: %s", task_id, e) + return {"success": False, "error": f"Failed to create browser session: {str(e)}"} + # Cleanup stops the supervisor before closing the backend; keep it stopped. + if command != "close" and session_info.get("cdp_url"): + _bt._ensure_cdp_supervisor(task_id) + + # Build the command with the appropriate backend flag. + # Cloud mode: --cdp connects to Browserbase. + # Local mode: --session launches a local headless Chromium. + # The rest of the command (--json, command, args) is identical. + if session_info.get("cdp_url"): + # Cloud mode — connect to remote Browserbase browser via CDP + # IMPORTANT: Do NOT use --session with --cdp. In agent-browser >=0.13, + # --session creates a local browser instance and silently ignores --cdp. + backend_args = ["--cdp", session_info["cdp_url"]] + else: + # Local mode — launch Chromium (headless by default, headed when configured) + backend_args = ["--session", session_info["session_name"]] + if _bt._is_headed_mode(): + backend_args.append("--headed") + + # Lightpanda engine injection (local mode only, agent-browser v0.25.3+). + # Use the resolved session backend rather than global cloud-provider state: + # hybrid private-URL routing can create a local sidecar while a cloud + # provider remains configured for public URLs. + engine = _engine_override or _bt._get_browser_engine() + if engine != "auto" and not _bt._is_camofox_mode() and not session_info.get("cdp_url"): + backend_args += ["--engine", engine] + + cmd_parts = _bt._agent_browser_argv(browser_cmd) + backend_args + ["--json", command] + args + + try: + task_socket_dir = _bt._prepare_session_socket_dir(session_info["session_name"]) + _bt.logger.debug("browser cmd=%s task=%s socket_dir=%s (%d chars)", + command, task_id, task_socket_dir, len(task_socket_dir)) + browser_env = _bt._agent_browser_command_env(task_socket_dir) + + # Chromium-only launch flags are rejected by Lightpanda. Strip both + # the current and legacy variables for Lightpanda commands; explicit + # Chrome commands and fallback use the shared Chromium policy. + if engine == "lightpanda": + _stripped_args = browser_env.pop("AGENT_BROWSER_ARGS", None) + _stripped_flags = browser_env.pop("AGENT_BROWSER_CHROME_FLAGS", None) + if _stripped_args is not None or _stripped_flags is not None: + _bt.logger.debug( + "browser: stripped Chromium-only AGENT_BROWSER_ARGS/" + "AGENT_BROWSER_CHROME_FLAGS for Lightpanda command %s " + "(agent-browser rejects them with --engine lightpanda)", + command, + ) + else: + _bt._apply_chromium_sandbox_args(browser_env) + + stdout_path = os.path.join(task_socket_dir, f"_stdout_{command}") + stderr_path = os.path.join(task_socket_dir, f"_stderr_{command}") + proc = _bt._popen_agent_browser(cmd_parts, browser_env, task_socket_dir, command) + + try: + proc.wait(timeout=timeout) + except subprocess.TimeoutExpired: + proc.kill() + proc.wait() + stdout, stderr = _bt._read_command_output_files(stdout_path, stderr_path) + _bt._unlink_command_output_files(stdout_path, stderr_path) + _bt._handle_browser_command_timeout(task_id, session_info, task_socket_dir) + if stderr and stderr.strip(): + _bt.logger.warning( + "browser '%s' stderr after timeout: %s", + command, + stderr.strip()[:500], + ) + _bt.logger.warning("browser '%s' timed out after %ds (task=%s, socket_dir=%s)", + command, timeout, task_id, task_socket_dir) + result = { + "success": False, + "error": _bt._format_browser_timeout_error(command, timeout, stdout, stderr), + } + # Fall through to fallback check below + else: + with open(stdout_path, "r", encoding="utf-8") as f: + stdout = f.read() + with open(stderr_path, "r", encoding="utf-8") as f: + stderr = f.read() + _bt._unlink_command_output_files(stdout_path, stderr_path) + result = _bt._interpret_browser_command_output(command, stdout, stderr, proc.returncode) + + except Exception as e: + _bt.logger.warning("browser '%s' exception: %s", command, e, exc_info=True) + result = {"success": False, "error": str(e)} + + # --- Lightpanda automatic Chrome fallback --- + # If engine is lightpanda and the result looks broken, retry with Chrome. + # This runs for ALL exit paths (timeout, empty, non-JSON, nonzero rc, parsed). + fallback_reason = _bt._lightpanda_fallback_reason(engine, command, result) + if fallback_reason: + _bt.logger.info( + "Lightpanda fallback: retrying '%s' with Chrome (task=%s): %s", + command, + task_id, + fallback_reason, + ) + # For screenshots, use the dedicated Chrome fallback helper + # (spins up a separate Chrome session to the same URL). + if command == "screenshot": + fallback_result = _bt._chrome_fallback_screenshot(task_id, args or [], timeout) + else: + fallback_result = _bt._run_chrome_fallback_command(task_id, command, args, timeout) + return _bt._annotate_lightpanda_fallback(fallback_result, fallback_reason) + + return result diff --git a/tools/browser_tool_vision.py b/tools/browser_tool_vision.py new file mode 100644 index 0000000000..b7e2ef90da --- /dev/null +++ b/tools/browser_tool_vision.py @@ -0,0 +1,181 @@ +"""browser_vision helpers: Lightpanda pre-route, native provider vision, auxiliary-LLM screenshot analysis. + +Split out of ``tools/browser_tool.py``; every name is re-imported there so +``tools.browser_tool.`` keeps resolving (and monkeypatching). Origin +symbols and module state are read/written through ``_bt`` (the origin module, +resolved per call by :func:`tools.browser_tool_origin.origin_module`) so +``patch("tools.browser_tool.X")`` is honoured and no import cycle exists. +""" + +import os +import shutil +from pathlib import Path +from typing import Any, Dict, Optional, Tuple + +from tools.browser_tool_origin import origin_module as _origin + + +def _vision_mode_label() -> str: + _bt = _origin() + _cp = _bt._get_cloud_provider() + return "local" if _cp is None else f"cloud ({_cp.provider_name()})" + + +def _lightpanda_vision_preroute( + effective_task_id: str, annotate: bool, screenshot_path: Path, +) -> Tuple[bool, Optional[str], Path]: + """Capture the vision screenshot through the Chrome fallback when Lightpanda is the engine. + + Lightpanda has no graphical renderer, so the normal path would fail with a + CDP error or return a placeholder PNG. Returns ``(prerouted, fallback_warning, + screenshot_path)``; on fallback failure ``prerouted`` is False and the caller + takes the normal screenshot path (forcing Chrome) so ``_run_browser_command`` + still produces the standard fallback metadata/error. + """ + _bt = _origin() + engine = _bt._get_browser_engine() + if engine != "lightpanda" or not _bt._should_inject_engine(engine): + return False, None, screenshot_path + _bt.logger.debug("browser_vision: pre-routing screenshot to Chrome (engine=lightpanda)") + screenshot_args = ["--annotate"] if annotate else [] + fb_result = _bt._chrome_fallback_screenshot(effective_task_id, screenshot_args, _bt._get_command_timeout()) + fb_result = _bt._annotate_lightpanda_fallback(fb_result, _bt._LP_VISION_FALLBACK_REASON) + if not fb_result.get("success"): + _bt.logger.warning("Lightpanda Chrome fallback vision screenshot failed: %s", fb_result.get("error")) + return False, None, screenshot_path + fb_path = fb_result.get("data", {}).get("path", "") + if fb_path and os.path.exists(fb_path): + import uuid as uuid_mod + from hermes_constants import get_hermes_dir + + screenshots_dir = get_hermes_dir("cache/screenshots", "browser_screenshots") + screenshots_dir.mkdir(parents=True, exist_ok=True) + persistent_path = screenshots_dir / f"browser_screenshot_{uuid_mod.uuid4().hex}.png" + shutil.copy2(fb_path, persistent_path) + screenshot_path = persistent_path + return True, fb_result.get("fallback_warning"), screenshot_path + + +def _native_vision_result( + screenshot_path: Path, question: str, annotate: bool, + result: Dict[str, Any], lp_fallback_warning: Optional[str], +) -> Dict[str, Any]: + """Multimodal tool-result envelope: the main model inspects the pixels itself. + + History-reuse cap: this embed is baked into the tool result and re-sent on + every later turn, exactly like vision_analyze's native path — apply the same + proactive resize so full-res screenshots can't enter immutable history + uncapped. The helper's stat/dimension quick-estimate skips the resize when + already under both caps; without Pillow it fails open to the raw bytes. + """ + from tools.vision_tools import ( + _EMBED_MAX_DIMENSION, + _EMBED_TARGET_BYTES, + _build_native_vision_tool_result, + _resize_image_for_vision, + ) + + data_url = _resize_image_for_vision( + screenshot_path, + mime_type="image/png", + max_base64_bytes=_EMBED_TARGET_BYTES, + max_dimension=_EMBED_MAX_DIMENSION, + force_jpeg=True, + ) + native_result = _build_native_vision_tool_result( + image_url=str(screenshot_path), + question=question, + image_data_url=data_url, + image_size_bytes=screenshot_path.stat().st_size, + ) + meta = native_result.setdefault("meta", {}) + meta["screenshot_path"] = str(screenshot_path) + if lp_fallback_warning: + meta["fallback_warning"] = lp_fallback_warning + if annotate and result.get("data", {}).get("annotations"): + meta["annotations"] = result["data"]["annotations"] + native_result["text_summary"] = ( + f"{native_result.get('text_summary', '')} " + f"Screenshot path: {screenshot_path}" + ).strip() + return native_result + + +def _analyze_screenshot_with_aux_llm(screenshot_path: Path, question: str) -> str: + """One-shot aux vision-LLM analysis (not baked into history), secret-redacted. + + Encodes at full resolution; on a size-related provider rejection the image + is downscaled once and retried. Timeout/temperature come from + ``auxiliary.vision.*`` — local vision models (llama.cpp, ollama) can take + well over 30s, so the default timeout is generous. + """ + _bt = _origin() + import base64 + + vision_prompt = ( + f"You are analyzing a screenshot of a web browser.\n\n" + f"User's question: {question}\n\n" + f"Provide a detailed and helpful answer based on what you see in the screenshot. " + f"If there are interactive elements, describe them. If there are verification challenges " + f"or CAPTCHAs, describe what type they are and what action might be needed. " + f"Focus on answering the user's specific question." + ) + _screenshot_bytes = screenshot_path.read_bytes() + _screenshot_b64 = base64.b64encode(_screenshot_bytes).decode("ascii") + data_url = f"data:image/png;base64,{_screenshot_b64}" + vision_model = _bt._get_vision_model() + _bt.logger.debug("browser_vision: analysing screenshot (%d bytes)", + len(_screenshot_bytes)) + + vision_timeout = 120.0 + vision_temperature = 0.1 + try: + from hermes_cli.config import load_config + _vision_cfg = _bt.cfg_get(load_config(), "auxiliary", "vision", default={}) + _vt = _vision_cfg.get("timeout") + if _vt is not None: + vision_timeout = float(_vt) + _vtemp = _vision_cfg.get("temperature") + if _vtemp is not None: + vision_temperature = float(_vtemp) + except Exception: + pass + + call_kwargs = { + "task": "vision", + "messages": [ + { + "role": "user", + "content": [ + {"type": "text", "text": vision_prompt}, + {"type": "image_url", "image_url": {"url": data_url}}, + ], + } + ], + "temperature": vision_temperature, + "timeout": vision_timeout, + } + if vision_model: + call_kwargs["model"] = vision_model + try: + response = _bt._lazy_call_llm(**call_kwargs) + except Exception as _api_err: + from tools.vision_tools import ( + _is_image_size_error, _resize_image_for_vision, _RESIZE_TARGET_BYTES, + ) + if not (_is_image_size_error(_api_err) and len(data_url) > _RESIZE_TARGET_BYTES): + raise + _bt.logger.info( + "Vision API rejected screenshot (%.1f MB); " + "auto-resizing to ~%.0f MB and retrying...", + len(data_url) / (1024 * 1024), + _RESIZE_TARGET_BYTES / (1024 * 1024), + ) + data_url = _resize_image_for_vision(screenshot_path, mime_type="image/png") + call_kwargs["messages"][0]["content"][1]["image_url"]["url"] = data_url + response = _bt._lazy_call_llm(**call_kwargs) + + analysis = (response.choices[0].message.content or "").strip() + # Redact secrets the vision LLM may have read from the screenshot. + from agent.redact import redact_sensitive_text + return redact_sensitive_text(analysis)