refactor(tools): unify never-raise/liveness helpers and tighten supervisor frame/dialog paths
This commit is contained in:
@@ -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: '<peer>/<agent>' 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",
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
``<handle|profile>@<connection-id>``; ``"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
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user