From 4025afc45a169a9d6578efdd43878415cb829a18 Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Wed, 2 Sep 2026 22:47:50 -0700 Subject: [PATCH] refactor(tools): unify never-raise/liveness helpers and tighten supervisor frame/dialog paths --- tools/bot_mode_dm.py | 67 ++++++++++--------------- tools/bot_mode_probe.py | 77 +++++++++++------------------ tools/bot_relay.py | 37 +++++--------- tools/browser_extension_router.py | 4 +- tools/browser_supervisor.py | 49 +++++++----------- tools/browser_supervisor_dialogs.py | 71 ++++++++++---------------- tools/browser_supervisor_frames.py | 48 +++++++----------- 7 files changed, 130 insertions(+), 223 deletions(-) diff --git a/tools/bot_mode_dm.py b/tools/bot_mode_dm.py index e0e3b9f31f..ebad41e988 100644 --- a/tools/bot_mode_dm.py +++ b/tools/bot_mode_dm.py @@ -29,21 +29,18 @@ import time from pathlib import Path from typing import Any, Optional -# Top-level imports stay stdlib-only: this module is also executed directly as the -# background delivery runner (``python bot_mode_dm.py --run-delivery …``), where -# ``tools.*`` is resolved from whichever install is on sys.path. Hermes-side -# helpers are imported lazily inside the functions that need them. +# Top-level imports stay stdlib-only: this module also runs directly as the background +# delivery runner (``python bot_mode_dm.py --run-delivery …``); Hermes helpers import lazily. logger = logging.getLogger(__name__) MESSAGE_AGENT_TOOL_NAME = "message_agent" -# Message body cap — generous for real work products, small enough that a -# runaway paste can't turn one DM into a context bomb on the recipient. +# Message body cap — generous for real work, small enough that a runaway paste can't +# turn one DM into a context bomb on the recipient. MESSAGE_MAX_CHARS = 16000 - -# A runner normally owns and removes each DM file. This bounds the residual -# plaintext lifetime if the machine dies between spawn ack and the runner's finally. +# A runner owns and removes each DM file; this bounds residual plaintext lifetime if +# the machine dies between spawn ack and the runner's finally. _DM_DIR_NAME = "hermes-dm" _DM_STALE_SECONDS = 24 * 60 * 60 @@ -181,8 +178,7 @@ def message_agent_tool(target: str = "", message: str = "", task_id: Optional[st except Exception as exc: # pragma: no cover — defensive return _err(f"Bot Mode gate check failed: {exc}") - root = _hermes_root(Path(home)) - me = _self_profile_name(Path(home)) + root, me = _hermes_root(Path(home)), _self_profile_name(Path(home)) roster = [name for name, _dir in _roster(root)] peers = _peers(root) teammates = [_handle(n) for n in roster if n != me] @@ -200,26 +196,21 @@ def message_agent_tool(target: str = "", message: str = "", task_id: Optional[st raw_target = str(target or "").strip().lstrip("@") if not raw_target: return _roster_err("target is required.") - - sender_handle = _handle(me) - content = f"Message from 🤖 {sender_handle} (@{sender_handle}): " + body + content = f"Message from 🤖 {_handle(me)} (@{_handle(me)}): " + body + delivery = dict(task_id=task_id, agent=agent) # Peer target: '/' or a bare registered peer name. peer_match = _PEER_TARGET_RE.match(raw_target) - bare_peer = raw_target.lower() if raw_target.lower() in peers else None - if peer_match or bare_peer: - peer_name = peer_match.group(1) if peer_match else bare_peer - peer_profile = peer_match.group(2) if peer_match else None + if peer_match or raw_target.lower() in peers: + peer_name, peer_profile = peer_match.groups() if peer_match else (raw_target.lower(), None) if peer_name not in peers: return _roster_err(f"No registered peer named '{peer_name}'.") dm_target = f"{peer_name}/{peer_profile}" if peer_profile else peer_name - # Pin the registry-owning profile: `hermes peer` resolves bot_peers via - # the profile-scoped load_config(), while the roster above reads the - # machine-root config — the CLI must run in that same profile or a - # secondary-profile bot sees an empty registry ("No peer named"). + # Pin the registry-owning profile: `hermes peer` resolves bot_peers via the profile-scoped + # load_config(), while the roster above reads the machine-root config — the CLI must run + # in that same profile or a secondary-profile bot sees an empty registry. return _start_delivery(["hermes", "-p", _self_profile_name(root), "peer", "dm", dm_target], content, - f"@{peer_profile or peer_name} on peer '{peer_name}'", - stdin_file=True, task_id=task_id, agent=agent) + f"@{peer_profile or peer_name} on peer '{peer_name}'", stdin_file=True, **delivery) # Local teammate. is_local_shape = bool(_LOCAL_TARGET_RE.match(raw_target)) @@ -227,11 +218,10 @@ def message_agent_tool(target: str = "", message: str = "", task_id: Optional[st return _roster_err(f"Invalid target: {raw_target!r}.") resolved = _resolve_local_name(raw_target, roster) if is_local_shape else None if resolved is None or resolved == me: - # Unknown locally, or same-name target on ANOTHER connection (this - # gateway's 'default' messaging the cloud 'default'): every gateway - # connected to the user's Desktop is reachable via the relay roster, so - # try that before reporting a resolution failure / self-message. - relayed = _try_relay_delivery(root, raw_target, content, me, task_id=task_id, agent=agent) + # Unknown locally, or same-name target on ANOTHER connection (this gateway's 'default' + # messaging the cloud 'default'): every Desktop-connected gateway is reachable via the + # relay roster, so try that before reporting a resolution failure / self-message. + relayed = _try_relay_delivery(root, raw_target, content, me, **delivery) if relayed is not None: return relayed if resolved == me: @@ -240,7 +230,7 @@ def message_agent_tool(target: str = "", message: str = "", task_id: Optional[st "machine, or on a registered peer. Pick a name from the roster " "(roles are listed in your system prompt).") return _start_delivery(["hermes", "-p", resolved, *BOT_CHAT_TURN_ARGS], content, f"@{_handle(resolved)}", - stdin_file=False, task_id=task_id, agent=agent) + stdin_file=False, **delivery) def _try_relay_delivery(root: Path, raw_target: str, content: str, me: str, *, @@ -278,8 +268,7 @@ def _try_relay_delivery(root: Path, raw_target: str, content: str, me: str, *, def _dm_dir() -> Path: uid_getter = getattr(os, "getuid", None) uid = uid_getter() if callable(uid_getter) else None - dirname = f"{_DM_DIR_NAME}-{uid}" if uid is not None else _DM_DIR_NAME - path = Path(tempfile.gettempdir()) / dirname + path = Path(tempfile.gettempdir()) / (f"{_DM_DIR_NAME}-{uid}" if uid is not None else _DM_DIR_NAME) path.mkdir(mode=0o700, exist_ok=True) # Shared POSIX temp roots need a per-user directory. Fail closed if an # attacker pre-created the expected path or replaced it with a symlink. @@ -320,8 +309,7 @@ def _write_dm_file(content: str) -> str: with os.fdopen(fd, "w", encoding="utf-8") as f: f.write(content) except BaseException: - # fdopen owns the descriptor once it succeeds, but if fdopen itself - # failed the raw descriptor is still ours. Closing twice is harmless. + # If fdopen itself failed the raw descriptor is still ours; closing twice is harmless. with contextlib.suppress(OSError): os.close(fd) _unlink_dm_file(path) @@ -359,8 +347,7 @@ def _run_local_turn(argv: list[str], dm_file: str) -> int: if proc.returncode != 0: from tools.bot_failure_reasons import RETRY_NONE, classify_agent_error, retry_action - detail = (proc.stderr or proc.stdout or "").strip()[-500:] - if retry_action(classify_agent_error(detail)) != RETRY_NONE: + if retry_action(classify_agent_error((proc.stderr or proc.stdout or "").strip()[-500:])) != RETRY_NONE: proc = _turn() if proc.returncode != 0 and "already has a live owner" in (proc.stderr or ""): # The target's Bot Chat is held live by another surface (Desktop); the turn @@ -402,9 +389,8 @@ def _delivery_command(argv: list[str], dm_file: str, *, stdin_file: bool) -> str runner_argv = [sys.executable, str(Path(__file__).resolve()), "--run-delivery", "stdin" if stdin_file else "query-file", dm_file, *argv] if sys.platform == "win32": - # The tracked local backend uses Git Bash on native Windows: forward - # slashes keep native drive paths executable there; backslash paths are - # parsed as command names and die with exit 127 before the runner starts. + # The tracked local backend uses Git Bash on native Windows: forward slashes keep drive + # paths executable there; backslash paths are parsed as command names (exit 127). runner_argv = [part.replace("\\", "/") for part in runner_argv] return shlex.join(runner_argv) @@ -441,8 +427,7 @@ def _spawn_delivery(command: str, label: str, *, dm_file: Optional[str] = None, return _err(f"Delivery to {label} failed to start: {parsed['error']}") if not proc_id: return _err(f"Delivery to {label} failed to start: no process id returned") - # From here the background runner owns the file and removes it only - # after the local query-file or peer stdin consumer has finished. + # From here the background runner owns the file (removed after the consumer finishes). transferred = True return json.dumps({ "status": "sent", diff --git a/tools/bot_mode_probe.py b/tools/bot_mode_probe.py index d19cbed2f7..4b5dcd35b2 100644 --- a/tools/bot_mode_probe.py +++ b/tools/bot_mode_probe.py @@ -1,14 +1,12 @@ """Bot Mode roster probe — canonical Bot Chat system prompt section. -When any profile on this install carries ``ui_meta['hermes-bots']`` in its -profile.yaml (Bot-Mode-managed), a bot's canonical "Bot Chat" session — and ONLY -that session (agent/system_prompt.py enforces the ``BOT_CHAT_TITLE`` gate) — gets -a "Messaging other agents" section. Silent (``""``) when no profile is managed, -when the profile's SOUL.md already carries the heading (legacy plugin-appended -text must never double up), or on any error — a prompt build must never crash. -Cached per (process, home) so compression-triggered rebuilds produce identical -bytes. Toggle: ``agent.bot_mode_protocol`` (default True). Also hosts the -path/roster helpers shared by ``bot_mode_dm`` and ``bot_relay``. +When any profile carries ``ui_meta['hermes-bots']`` in profile.yaml (Bot-Mode-managed), +a bot's canonical "Bot Chat" session — ONLY that session (agent/system_prompt.py enforces +the ``BOT_CHAT_TITLE`` gate) — gets a "Messaging other agents" section. Silent (``""``) +when no profile is managed, when SOUL.md already carries the heading (legacy plugin text +must never double up), or on any error. Cached per (process, home) so compression rebuilds +produce identical bytes. Toggle: ``agent.bot_mode_protocol``. Also hosts path/roster +helpers shared by ``bot_mode_dm`` and ``bot_relay``. """ from __future__ import annotations @@ -69,11 +67,8 @@ def _roster(root: Path) -> list[tuple[str, Path]]: def _read_yaml_dict(path: Path, needle: str | None = None) -> dict | None: - """YAML mapping at ``path``, or None when missing / not a mapping / unreadable. - - ``needle``: cheap substring precheck that skips the YAML parse on the - dominant (unmanaged) path — the key is absent from most installs. - """ + """YAML mapping at ``path``, or None when missing / not a mapping / unreadable. ``needle``: + cheap substring precheck that skips the YAML parse on the dominant (unmanaged) path.""" def _load(): if not path.is_file(): return None @@ -104,12 +99,9 @@ def _any_managed(root: Path) -> bool: def is_bot_mode_managed(home: str | os.PathLike | None = None) -> bool: - """True when ANY profile on this install is Bot-Mode-managed. Never raises. - - The tool-injection gate for ``message_agent`` — deliberately independent of - the protocol section's emptiness: a SOUL.md carrying the legacy protocol - gets an empty section but must still get the tool. - """ + """True when ANY profile on this install is Bot-Mode-managed. Never raises. The + ``message_agent`` injection gate — deliberately independent of the protocol section's + emptiness: a SOUL.md carrying the legacy protocol gets an empty section but still gets the tool.""" return _swallow(lambda: _any_managed(_hermes_root(_resolve_home(home))), False) @@ -236,12 +228,9 @@ def _build_section(home: Path) -> str: def get_bot_mode_protocol_section(home: str | os.PathLike | None = None, *, force_refresh: bool = False) -> str: - """Cached probe entry point — one filesystem pass per (process, home). - - ``home`` should be the AGENT'S OWN resolved home (session-db derived), not - the ambient HERMES_HOME — build threads can lose the ContextVar override - and the env var would then name the wrong profile. - """ + """Cached probe entry point — one filesystem pass per (process, home). ``home`` should be + the AGENT'S OWN resolved home (session-db derived), not ambient HERMES_HOME — build threads + can lose the ContextVar override and the env var would then name the wrong profile.""" resolved = str(_resolve_home(home)) with _lock: if force_refresh or resolved not in _cached: @@ -250,26 +239,20 @@ def get_bot_mode_protocol_section(home: str | os.PathLike | None = None, *, forc # ── capability epoch ───────────────────────────────────────────────────────── -# -# Bot Chat sessions are effectively eternal, so "build the prompt once" would -# strand capability changes (skills, toolsets, MCP, SOUL, roster, peers) forever. -# The fingerprint hashes exactly that surface; the built Bot Chat prompt embeds -# it and agent/conversation_loop.py rebuilds only when the stored epoch differs -# from disk — a loud, once-per-change cache break, never per-turn drift. +# Bot Chat sessions are effectively eternal, so "build the prompt once" would strand +# capability changes (skills, toolsets, MCP, SOUL, roster, peers) forever. The fingerprint +# hashes exactly that surface; the built prompt embeds it and agent/conversation_loop.py +# rebuilds only when the stored epoch differs from disk — once per change, never per-turn drift. _EPOCH_PREFIX = "Capability epoch: " _EPOCH_RE_TEXT = r"Capability epoch: ([0-9a-f]{12})" def capability_fingerprint(home: str | os.PathLike | None = None) -> str: - """12-hex digest of the capability surface for ``home``'s profile. - - Sources: disabled skills + enabled toolsets + MCP config (config.yaml), - SOUL.md bytes, installed skill names, the Bot-Mode roster (+ roles), peers - and the relay roster. Deliberately NOT cached — the point is detecting - on-disk drift against the epoch embedded in a stored prompt. Never raises - ("unavailable" on failure). - """ + """12-hex digest of the capability surface for ``home``'s profile: disabled skills + + enabled toolsets + MCP config, SOUL.md bytes, installed skill names, the Bot-Mode roster + (+ roles), peers and the relay roster. Deliberately NOT cached — the point is detecting + on-disk drift against a stored prompt's epoch. Never raises ("unavailable" on failure).""" import hashlib import json @@ -350,15 +333,11 @@ def stored_prompt_capability_stale(stored_prompt: str, home: str | os.PathLike | def stored_bot_chat_prompt_needs_upgrade(stored_prompt: str, home: str | os.PathLike | None = None) -> bool: - """True when a Bot Chat session's stored prompt PREDATES the epoch mechanism. - - Legacy prompts carry neither protocol section nor epoch stamp, so the - staleness check (stamped prompts only) would strand them forever. One-time - migration: the caller must only ask for sessions titled "Bot Chat", and we - rebuild only when the probe would actually emit a section — a SOUL.md that - already carries the legacy protocol yields an empty section, and rebuilding - would produce another unstamped prompt and loop. Fails closed to "no upgrade". - """ + """True when a Bot Chat session's stored prompt PREDATES the epoch mechanism. Legacy + prompts carry neither section nor stamp, so the staleness check (stamped only) would strand + them forever. The caller must only ask for sessions titled "Bot Chat"; we rebuild only when + the probe would actually emit a section — a SOUL.md carrying the legacy protocol yields an + empty section, and rebuilding would mint another unstamped prompt and loop. Fails closed.""" text = stored_prompt or "" if _EPOCH_PREFIX in text or _PROTOCOL_HEADING in text: return False diff --git a/tools/bot_relay.py b/tools/bot_relay.py index ff0ba85712..a714b8cce8 100644 --- a/tools/bot_relay.py +++ b/tools/bot_relay.py @@ -130,13 +130,11 @@ def write_remote_roster(root: Path | str, rows: Any) -> int: """Atomically persist the Desktop-pushed remote roster. Returns count.""" base = _ensure_dirs(root) by_key: dict[tuple[str, str], dict] = {} - for row in rows if isinstance(rows, list) else []: - norm = _normalize_roster_row(row) - if norm: - by_key.setdefault((norm["connection_id"], norm["profile"]), norm) + for norm in filter(None, map(_normalize_roster_row, rows if isinstance(rows, list) else [])): + by_key.setdefault((norm["connection_id"], norm["profile"]), norm) cleaned = [by_key[k] for k in sorted(by_key)] - payload = {"updated_at": int(time.time()), "agents": cleaned} - _atomic_write_json(base / ROSTER_FILE, payload, prefix=".roster-", sort_keys=True) + _atomic_write_json(base / ROSTER_FILE, {"updated_at": int(time.time()), "agents": cleaned}, + prefix=".roster-", sort_keys=True) return len(cleaned) @@ -145,9 +143,7 @@ def read_remote_roster(root: Path | str) -> list[dict]: try: data = json.loads((relay_root(root) / ROSTER_FILE).read_text(encoding="utf-8")) agents = data.get("agents") if isinstance(data, dict) else None - if not isinstance(agents, list): - return [] - return [r for r in (_normalize_roster_row(a) for a in agents) if r] + return [r for r in map(_normalize_roster_row, agents) if r] if isinstance(agents, list) else [] except FileNotFoundError: return [] except Exception: @@ -158,15 +154,9 @@ def read_remote_roster(root: Path | str) -> list[dict]: def resolve_remote_target(raw_target: str, roster: list[dict]) -> Any: """Matched row for a bare handle/profile (unique across connections) or ``@``; ``"ambiguous"`` for a bare form on several connections; None otherwise.""" - want = str(raw_target or "").strip().lstrip("@") - if not want: + want, at, conn = (p.strip() for p in str(raw_target or "").strip().lstrip("@").partition("@")) + if not want or (at and not conn): return None - conn: Optional[str] = None - if "@" in want: - want, _, conn = want.partition("@") - want, conn = want.strip(), conn.strip() - if not want or not conn: - return None matches = [row for row in roster if want.lower() in (row["handle"].lower(), row["profile"].lower()) and (not conn or row["connection_id"].lower() == conn.lower())] if not matches: @@ -198,17 +188,14 @@ def _target_liveness(root: Path | str, target: dict) -> Optional[bool]: age = time.time() - (relay_root(root) / ROSTER_FILE).stat().st_mtime except OSError: return None - if age > ROSTER_FRESH_SECONDS: - return None - roster = read_remote_roster(root) + roster = read_remote_roster(root) if age <= ROSTER_FRESH_SECONDS else [] if not roster: return None key = (str(target.get("connection_id") or ""), str(target.get("profile") or "")) - for row in roster: - if (row["connection_id"], row["profile"]) == key: - online = row.get("online") - return online if isinstance(online, bool) else None - return False # fresh roster no longer lists the target — offline + row = next((r for r in roster if (r["connection_id"], r["profile"]) == key), None) + if row is None: + return False # fresh roster no longer lists the target — offline + return row["online"] if isinstance(row.get("online"), bool) else None except Exception: logger.debug("bot_relay liveness check failed", exc_info=True) return None diff --git a/tools/browser_extension_router.py b/tools/browser_extension_router.py index 626fcc26b4..70e62f9418 100644 --- a/tools/browser_extension_router.py +++ b/tools/browser_extension_router.py @@ -52,8 +52,8 @@ def extension_controller_available(action: str) -> bool: if not browser_control_enabled(): return False - session_id, principal_id, transport_family = _bound_identity() - if not session_id or not principal_id or not transport_family: + session_id, principal_id, transport_family = identity = _bound_identity() + if not all(identity): return False broker = get_browser_control_broker() scope = broker.scope_for_session(session_id=session_id, principal_id=principal_id, transport_family=transport_family) diff --git a/tools/browser_supervisor.py b/tools/browser_supervisor.py index 67ab90d7e5..1b79ce7a3d 100644 --- a/tools/browser_supervisor.py +++ b/tools/browser_supervisor.py @@ -1,15 +1,11 @@ """Persistent CDP supervisor for browser dialog + frame detection. -One ``CDPSupervisor`` runs per Hermes ``task_id`` with a reachable CDP endpoint. -It holds one persistent WebSocket, subscribes to ``Page`` / ``Runtime`` / -``Target`` events on every attached session (top page + auto-attached OOPIF / -worker targets), and exposes pending dialogs + frame tree through a -thread-safe snapshot that tool handlers read synchronously. - -Not in the agent's tool schema: output reaches the agent via ``browser_snapshot`` -and ``browser_dialog``. Dialog capture lives in ``browser_supervisor_dialogs``, -frame tracking in ``browser_supervisor_frames``; both are mixed into -``CDPSupervisor`` and their public names are re-exported here. +One ``CDPSupervisor`` per Hermes ``task_id`` with a reachable CDP endpoint: one +persistent WebSocket, ``Page`` / ``Runtime`` / ``Target`` events on every attached +session (top page + auto-attached OOPIF / worker targets), pending dialogs + frame +tree exposed via a thread-safe snapshot. Not in the tool schema — output reaches the +agent via ``browser_snapshot`` / ``browser_dialog``. Dialog capture and frame tracking +are mixins (``browser_supervisor_dialogs`` / ``browser_supervisor_frames``) re-exported here. Design spec: ``website/docs/developer-guide/browser-supervisor.md``. """ @@ -43,12 +39,9 @@ logger = logging.getLogger(__name__) def _redact_cdp_error_text(exc: object) -> str: """Redact CDP endpoint credentials from an exception's (or URL's) string form. - - ``websockets`` bakes the raw target URL (``?token=`` / ``user:pass@``) into - its exception messages, so every egress point that turns such an exception - into log text or a re-raised message MUST route through here; falls back to - a fixed sentinel if redaction itself raises, erring toward masking. - """ + ``websockets`` bakes the raw URL (``?token=`` / ``user:pass@``) into its exception + messages, so every egress point turning one into log/re-raise text MUST route through + here; falls back to a fixed sentinel if redaction itself raises (err toward masking).""" try: from agent.redact import redact_cdp_url @@ -148,8 +141,7 @@ class CDPSupervisor(DialogSupervisionMixin, FrameTrackingMixin): self.stop() raise TimeoutError(f"CDP supervisor did not attach within {timeout}s " f"(cdp_url={_redact_cdp_error_text(self.cdp_url)[:80]}...)") - if self._start_error is not None: - err = self._start_error + if (err := self._start_error) is not None: self.stop() # ``err`` is a raw ``websockets`` exception embedding the full cdp_url (token / # userinfo): re-raise redacted, ``from None`` so the traceback chain leaks nothing. @@ -184,11 +176,8 @@ class CDPSupervisor(DialogSupervisionMixin, FrameTrackingMixin): def respond_to_dialog(self, action: str, *, prompt_text: Optional[str] = None, dialog_id: Optional[str] = None, timeout: float = 10.0) -> Dict[str, Any]: - """Accept/dismiss a pending dialog (sync bridge onto the supervisor loop). - - Returns ``{"ok": True, "dialog": {...}}`` or ``{"ok": False, "error": ...}`` - for recoverable errors (no dialog, ambiguous dialog_id, inactive). - """ + """Accept/dismiss a pending dialog (sync bridge onto the supervisor loop). Returns + ``{"ok": True, "dialog"}`` or ``{"ok": False, "error"}`` for recoverable errors.""" if action not in {"accept", "dismiss"}: return _fail(f"action must be 'accept' or 'dismiss', got {action!r}") with self._state_lock: @@ -227,9 +216,9 @@ class CDPSupervisor(DialogSupervisionMixin, FrameTrackingMixin): if loop is None or not loop.is_running(): return _fail("supervisor loop is not running") with self._state_lock: - if not self._active: - return _fail("supervisor is not active") - session_id = self._page_session_id + active, session_id = self._active, self._page_session_id + if not active: + return _fail("supervisor is not active") if not session_id: return _fail("supervisor has no attached page session") @@ -427,8 +416,7 @@ class CDPSupervisor(DialogSupervisionMixin, FrameTrackingMixin): class _SupervisorRegistry: - """Process-global (task_id → supervisor) map with idempotent start/stop. - One instance, exposed as ``SUPERVISOR_REGISTRY``; mutations go through ``_lock``.""" + """Process-global (task_id → supervisor) map with idempotent start/stop (``SUPERVISOR_REGISTRY``).""" def __init__(self) -> None: self._lock = threading.Lock() @@ -444,9 +432,8 @@ class _SupervisorRegistry: def get_or_start(self, task_id: str, cdp_url: str, *, dialog_policy: str = DEFAULT_DIALOG_POLICY, dialog_timeout_s: float = DEFAULT_DIALOG_TIMEOUT_S, start_timeout: float = 15.0) -> CDPSupervisor: - """Idempotently ensure a supervisor is running for ``(task_id, cdp_url)``. - An existing supervisor bound to a different ``cdp_url`` (or unhealthy: - dead thread / stopped loop) is stopped and replaced.""" + """Idempotently ensure a supervisor runs for ``(task_id, cdp_url)``; one bound to a + different ``cdp_url`` or unhealthy (dead thread / stopped loop) is stopped and replaced.""" with self._lock: existing = self._by_task.get(task_id) if existing is not None: diff --git a/tools/browser_supervisor_dialogs.py b/tools/browser_supervisor_dialogs.py index 4150610e17..cc6a8ad56a 100644 --- a/tools/browser_supervisor_dialogs.py +++ b/tools/browser_supervisor_dialogs.py @@ -1,15 +1,11 @@ """Dialog capture + response half of the CDP supervisor. -Two capture paths feed the same ``PendingDialog`` queue: native -``Page.javascriptDialogOpening`` events (answered with -``Page.handleJavaScriptDialog``), and the injected *dialog bridge* — a page -script that rewrites alert/confirm/prompt into a sync XHR to a magic host we -intercept via the CDP ``Fetch`` domain and answer with ``Fetch.fulfillRequest``. -The bridge works on Browserbase, whose CDP proxy auto-dismisses real native -dialogs, because the native dialog never fires. - -``DialogSupervisionMixin`` relies on state ``CDPSupervisor.__init__`` sets -(``_state_lock``, ``_pending_dialogs``, ``_recent_dialogs``, ``_dialog_watchdogs``, ``_cdp`` ...). +Two capture paths feed one ``PendingDialog`` queue: native ``Page.javascriptDialogOpening`` +events (answered with ``Page.handleJavaScriptDialog``), and the injected *dialog bridge* — +a page script rewriting alert/confirm/prompt into a sync XHR to a magic host we intercept +via the CDP ``Fetch`` domain and answer with ``Fetch.fulfillRequest``. The bridge works on +Browserbase (whose CDP proxy auto-dismisses native dialogs) because the native dialog never fires. +``DialogSupervisionMixin`` relies on state ``CDPSupervisor.__init__`` sets. """ from __future__ import annotations @@ -56,22 +52,19 @@ _VALID_POLICIES = frozenset( DEFAULT_DIALOG_POLICY = DIALOG_POLICY_MUST_RESPOND DEFAULT_DIALOG_TIMEOUT_S = 300.0 -# Last N closed dialogs kept in ``recent_dialogs`` so agents on backends that -# auto-dismiss server-side (Browserbase) can still observe that a dialog fired. +# Last N closed dialogs kept so agents on backends that auto-dismiss server-side +# (Browserbase) can still observe that a dialog fired. RECENT_DIALOGS_MAX = 20 -# Magic host the injected dialog bridge XHRs to. Intercepted via the CDP Fetch -# domain before any network resolution, so it never has to exist. Keep ASCII + -# URL-safe; Fetch patterns are gated on it. +# Magic host the bridge XHRs to; intercepted via CDP Fetch before any network +# resolution, so it never has to exist. Keep ASCII + URL-safe (Fetch patterns gate on it). DIALOG_BRIDGE_HOST = "hermes-dialog-bridge.invalid" DIALOG_BRIDGE_URL_PATTERN = f"http://{DIALOG_BRIDGE_HOST}/*" -# Injected into every frame via Page.addScriptToEvaluateOnNewDocument. Uses a -# sync GET with query params so the Fetch interceptor never parses a body; if -# the bridge is unreachable it returns null so the page still sees *some* -# behavior (the backend auto-dismisses). onbeforeunload is left native — it -# can't be prompted synchronously without racing navigation; the native-dialog -# fallback path still surfaces it in recent_dialogs. +# Injected into every frame via Page.addScriptToEvaluateOnNewDocument. Sync GET with +# query params so the Fetch interceptor never parses a body; unreachable bridge → null +# so the page still sees *some* behavior. onbeforeunload is left native (can't be +# prompted synchronously without racing navigation); the native path still records it. _DIALOG_BRIDGE_SCRIPT = r""" (() => { if (window.__hermesDialogBridgeInstalled) return; @@ -125,8 +118,7 @@ class PendingDialog: opened_at: float cdp_session_id: str # which attached CDP session the dialog fired in frame_id: Optional[str] = None - # Set when captured via the bridge XHR path: respond via Fetch.fulfillRequest, - # NOT Page.handleJavaScriptDialog — the native dialog never fired. + # Bridge XHR path: respond via Fetch.fulfillRequest, NOT Page.handleJavaScriptDialog. bridge_request_id: Optional[str] = None def to_dict(self) -> Dict[str, Any]: @@ -161,13 +153,10 @@ class DialogSupervisionMixin: logger.debug("%s failed (%s): %s", method, what, e) async def _install_dialog_bridge(self, session_id: str) -> None: - """Install the dialog-bridge init script + Fetch interceptor on a session. - - Idempotent at the CDP level (Chromium de-dupes identical add-script - calls; Fetch.enable replaces prior patterns). The final Runtime.evaluate - injects into the already-loaded document so existing pages pick up the - override on reconnect. - """ + """Install the dialog-bridge init script + Fetch interceptor on a session. Idempotent at + the CDP level (Chromium de-dupes identical add-script calls; Fetch.enable replaces prior + patterns); the final Runtime.evaluate injects into the already-loaded document so + existing pages pick up the override on reconnect.""" sid = (session_id or "")[:16] steps = ( ("Page.addScriptToEvaluateOnNewDocument", {"source": _DIALOG_BRIDGE_SCRIPT, "runImmediately": True}, @@ -189,12 +178,9 @@ class DialogSupervisionMixin: ) async def _on_fetch_paused(self, params: Dict[str, Any], session_id: Optional[str]) -> None: - """Bridge XHR captured mid-flight — materialize as a pending dialog. - - The page's JS thread is blocked on the XHR until we Fetch.fulfillRequest - (from ``respond_to_dialog`` or the watchdog). Requests for other hosts - are forwarded unchanged so the page sees its own request. - """ + """Bridge XHR captured mid-flight — materialize as a pending dialog. The page's JS + thread is blocked on the XHR until we Fetch.fulfillRequest (agent or watchdog); + requests for other hosts are forwarded unchanged.""" url = str(params.get("request", {}).get("url") or "") request_id = params.get("requestId") if not request_id: @@ -212,10 +198,8 @@ class DialogSupervisionMixin: def _admit_dialog(self, *, type: str, message: str, default_prompt: str, session_id: Optional[str], frame_id: Optional[str], bridge_request_id: Optional[str] = None) -> None: """Create the dialog and apply the policy: auto-respond, or queue + arm the watchdog. - - Auto policies archive FIRST (tagged ``auto_policy``) so the ``closed`` - event that follows our own response isn't re-archived as ``remote``. - """ + Auto policies archive FIRST (tagged ``auto_policy``) so the ``closed`` event that + follows our own response isn't re-archived as ``remote``.""" self._dialog_seq += 1 dialog = PendingDialog( id=f"d-{self._dialog_seq}", type=type, message=message, default_prompt=default_prompt, @@ -306,10 +290,9 @@ class DialogSupervisionMixin: self._recent_dialogs = _trim_ring([*self._recent_dialogs, record], RECENT_DIALOGS_MAX) async def _on_dialog_closed(self, params: Dict[str, Any], session_id: Optional[str]) -> None: - # ``Page.javascriptDialogClosed`` carries only ``result``/``userInput``, not - # the message. Match by session id and clear the oldest native dialog on - # it — the JS thread blocks while a dialog is up, so at most one is in - # flight per session. Bridge dialogs resolve via Fetch.fulfillRequest. + # ``Page.javascriptDialogClosed`` carries only ``result``/``userInput``: match by + # session id and clear the oldest native dialog on it (the JS thread blocks while + # a dialog is up, so at most one is in flight). Bridge dialogs resolve via Fetch. with self._state_lock: candidate = next((d.id for d in self._pending_dialogs.values() if d.cdp_session_id == session_id and d.bridge_request_id is None), None) diff --git a/tools/browser_supervisor_frames.py b/tools/browser_supervisor_frames.py index 5832b8b0fc..b3d33fa8c8 100644 --- a/tools/browser_supervisor_frames.py +++ b/tools/browser_supervisor_frames.py @@ -39,10 +39,8 @@ class FrameInfo: def to_dict(self) -> Dict[str, Any]: d = {"frame_id": self.frame_id, "url": self.url, "origin": self.origin, "is_oopif": self.is_oopif} - for key, value in (("session_id", self.cdp_session_id), ("parent_frame_id", self.parent_frame_id), - ("name", self.name)): - if value: - d[key] = value + optional = (("session_id", self.cdp_session_id), ("parent_frame_id", self.parent_frame_id), ("name", self.name)) + d.update({k: v for k, v in optional if v}) return d @@ -71,33 +69,26 @@ class FrameTrackingMixin: if not frame_id: return with self._state_lock: - old = self._frames.get(frame_id) + old = self._frames.get(frame_id) or FrameInfo(frame_id, "", "", None, False, session_id) self._frames[frame_id] = FrameInfo( frame_id=frame_id, url=str(frame.get("url") or ""), origin=str(frame.get("securityOrigin") or frame.get("origin") or ""), - parent_frame_id=frame.get("parentId") or (old.parent_frame_id if old else None), - is_oopif=bool(old.is_oopif if old else False), - cdp_session_id=old.cdp_session_id if old else session_id, - name=str(frame.get("name") or (old.name if old else "")), + parent_frame_id=frame.get("parentId") or old.parent_frame_id, is_oopif=old.is_oopif, + cdp_session_id=old.cdp_session_id, name=str(frame.get("name") or old.name), ) def _on_frame_detached(self, params: Dict[str, Any], session_id: Optional[str]) -> None: - """Drop a frame only when it's truly gone. - - ``reason="swap"`` means the frame is migrating processes (e.g. promoted - to an OOPIF) — dropping it would hide the iframe. Even with ``remove`` - the parent only knows the child left ITS process; if we hold a live - child session for that frame_id it is still alive, so keep it until - Target.detached + a later frameDetached clear it. - """ + """Drop a frame only when it's truly gone. ``reason="swap"`` = migrating processes + (e.g. promoted to an OOPIF) — dropping would hide the iframe. Even with ``remove`` + the parent only knows the child left ITS process; a live child session means it's + still alive, so keep it until Target.detached + a later frameDetached clear it.""" frame_id = params.get("frameId") if not frame_id or str(params.get("reason") or "remove").lower() == "swap": return with self._state_lock: old = self._frames.get(frame_id) - if old and old.is_oopif and old.cdp_session_id: - return - self._frames.pop(frame_id, None) + if not (old and old.is_oopif and old.cdp_session_id): + self._frames.pop(frame_id, None) async def _on_target_attached(self, params: Dict[str, Any], session_id: Optional[str] = None) -> None: info = params.get("targetInfo") or {} @@ -113,7 +104,7 @@ class FrameTrackingMixin: old = self._frames.get(target_id) self._frames[target_id] = FrameInfo( frame_id=target_id, url=str(info.get("url") or ""), origin="", is_oopif=True, cdp_session_id=sid, - parent_frame_id=(old.parent_frame_id if old else None), name=str(info.get("title") or (old.name if old else "")), + parent_frame_id=old.parent_frame_id if old else None, name=str(info.get("title") or (old.name if old else "")), ) # Enable child domains off-loop: awaiting the replies here would deadlock # because only the reader can resolve those Futures. @@ -128,20 +119,15 @@ class FrameTrackingMixin: await self._install_dialog_bridge(sid) def _on_target_detached(self, params: Dict[str, Any], session_id: Optional[str] = None) -> None: - """Clear the session binding of frames on a detached child session. - - Frames are deliberately NOT dropped: Browserbase fires transient detaches - during page transitions while the iframe is still visible. Clearing - ``cdp_session_id`` just stops stale routing; ``Page.frameDetached`` - cleans up if the iframe truly goes away. - """ + """Clear the session binding of frames on a detached child session. Frames are + deliberately NOT dropped: Browserbase fires transient detaches during page transitions + while the iframe is still visible; ``Page.frameDetached`` cleans up if it truly goes away.""" sid = params.get("sessionId") if not sid: return with self._state_lock: - for fid, frame in list(self._frames.items()): - if frame.cdp_session_id == sid: - self._frames[fid] = replace(frame, cdp_session_id=None) + self._frames.update({fid: replace(f, cdp_session_id=None) for fid, f in self._frames.items() + if f.cdp_session_id == sid}) def _build_frame_tree_locked(self) -> Dict[str, Any]: """Capped frame_tree payload (must hold state lock). Top frame = one with