diff --git a/plugins/google_meet/__init__.py b/plugins/google_meet/__init__.py index df401e1a68..48af371dc0 100644 --- a/plugins/google_meet/__init__.py +++ b/plugins/google_meet/__init__.py @@ -1,13 +1,9 @@ """google_meet plugin — let the agent join a Meet call, transcribe it, follow up. -v1: transcribe-only. Spawns a headless Chromium via Playwright, joins the Meet -URL, enables live captions, scrapes them into a transcript file. The agent then -has the transcript in its workspace and can do whatever followup work it needs -using its regular tools. - -v2 (not in this PR): realtime duplex audio so the agent can speak in the -meeting, via OpenAI Realtime / Gemini Live + BlackHole / PulseAudio null-sink. -``meet_say`` exists as a stub today so the tool surface is stable. +Spawns a headless Chromium via Playwright, joins the Meet URL, enables live +captions and scrapes them into a transcript file in the agent's workspace. +Realtime mode additionally lets the agent speak (OpenAI Realtime + a virtual +audio device); remote nodes let the bot run on another machine. Explicit-by-design: only joins ``https://meet.google.com/`` URLs explicitly passed in. No calendar scanning, no auto-dial, no consent announcement. @@ -22,16 +18,8 @@ from plugins.google_meet import process_manager as pm from plugins.google_meet.cli import register_cli as _register_meet_cli from plugins.google_meet.cli import meet_command as _meet_command from plugins.google_meet.tools import ( - MEET_JOIN_SCHEMA, - MEET_LEAVE_SCHEMA, - MEET_SAY_SCHEMA, - MEET_STATUS_SCHEMA, - MEET_TRANSCRIPT_SCHEMA, - check_meet_requirements, - handle_meet_join, - handle_meet_leave, - handle_meet_say, - handle_meet_status, + MEET_JOIN_SCHEMA, MEET_LEAVE_SCHEMA, MEET_SAY_SCHEMA, MEET_STATUS_SCHEMA, MEET_TRANSCRIPT_SCHEMA, + check_meet_requirements, handle_meet_join, handle_meet_leave, handle_meet_say, handle_meet_status, handle_meet_transcript, ) @@ -48,11 +36,9 @@ _TOOLS = ( def _on_session_end(**kwargs) -> None: - """Best-effort cleanup — if a meet bot is still running when the session - ends, leave the call so we don't orphan a headless Chromium. + """Leave a still-running call so we don't orphan a headless Chromium. - No-ops when nothing is active. Swallows all exceptions — session end must - not fail because the bot cleanup hit an edge case. + Swallows all exceptions — session end must never fail on bot cleanup. """ try: status = pm.status() @@ -63,31 +49,17 @@ def _on_session_end(**kwargs) -> None: def register(ctx) -> None: - """Register tools, CLI, and lifecycle hooks. - - Called once by the plugin loader when the plugin is enabled via - ``plugins.enabled`` in config.yaml. - """ - # Windows is not supported in v1 — audio routing for v2 doesn't have a - # tested path there and guest-join Chromium is flakier. Refuse to register - # rather than half-working. + """Register tools, CLI, and lifecycle hooks (called once by the plugin loader).""" + # Windows: no tested audio-routing path and flakier guest-join Chromium — + # refuse to register rather than half-work. system = platform.system().lower() if system not in {"linux", "darwin"}: - logger.info( - "google_meet plugin: platform=%s not supported (linux/macos only)", - system, - ) + logger.info("google_meet plugin: platform=%s not supported (linux/macos only)", system) return for name, schema, handler, emoji in _TOOLS: - ctx.register_tool( - name=name, - toolset="google_meet", - schema=schema, - handler=handler, - check_fn=check_meet_requirements, - emoji=emoji, - ) + ctx.register_tool(name=name, toolset="google_meet", schema=schema, handler=handler, + check_fn=check_meet_requirements, emoji=emoji) ctx.register_cli_command( name="meet", diff --git a/plugins/google_meet/_jsonfile.py b/plugins/google_meet/_jsonfile.py new file mode 100644 index 0000000000..24a36a1a9a --- /dev/null +++ b/plugins/google_meet/_jsonfile.py @@ -0,0 +1,34 @@ +"""Tiny JSON file helpers shared by the bot, process manager, node registry and node server.""" + +from __future__ import annotations + +import json +from pathlib import Path +from typing import Any, Optional + + +def read_json(path: Path) -> Optional[Any]: + """Parsed JSON from *path*, or None when missing/unreadable/malformed.""" + if not path.is_file(): + return None + try: + return json.loads(path.read_text(encoding="utf-8")) + except (OSError, ValueError): + return None + + +def write_json_atomic(path: Path, data: Any, mode: Optional[int] = None) -> None: + """Write ``json.dumps(data, indent=2)`` via a ``.json.tmp`` sibling + rename. + + *mode* (e.g. ``0o600``) is applied to the temp file before the rename so + the final file never exists with looser permissions. + """ + path.parent.mkdir(parents=True, exist_ok=True) + tmp = path.with_suffix(".json.tmp") + tmp.write_text(json.dumps(data, indent=2), encoding="utf-8") + if mode is not None: + try: + tmp.chmod(mode) + except (OSError, NotImplementedError): # best-effort on non-POSIX filesystems + pass + tmp.replace(path) diff --git a/plugins/google_meet/audio_bridge.py b/plugins/google_meet/audio_bridge.py index 8159c21c17..4919580aa9 100644 --- a/plugins/google_meet/audio_bridge.py +++ b/plugins/google_meet/audio_bridge.py @@ -1,20 +1,14 @@ """Virtual audio bridge for feeding generated speech into Chrome's mic. -v2 module. Provisions a platform-specific virtual audio device so the -Meet bot's Chromium instance can be pointed at an input source we -control. The OpenAI Realtime client writes PCM bytes into this device; -Chrome reads them as if they were coming from a microphone. +Provisions a platform-specific virtual audio device the Meet bot's Chromium +can use as its input; the OpenAI Realtime client writes PCM into it. -Linux (primary): uses pactl (PulseAudio) to create a null-sink plus a -virtual source whose master is the null-sink's monitor. Callers set -PULSE_SOURCE= in Chrome's env and pass the fake-mic flag. +Linux: pactl creates a null-sink plus a virtual source whose master is the +sink's monitor. Callers set ``PULSE_SOURCE=`` in Chrome's env. -macOS: requires BlackHole 2ch to be installed. This module only -verifies its presence and returns the device name; routing OS default -input is left to the user (or a future switchaudio-osx integration) to -avoid surprising the user's system audio state. - -Windows: not supported in v2. +macOS: only verifies BlackHole 2ch is installed and returns its name; the +default-input switch is left to the user so we never surprise their system +audio state. Windows: unsupported. """ from __future__ import annotations @@ -27,11 +21,21 @@ from typing import Optional _BLACKHOLE_DEVICE = "BlackHole 2ch" +def _pactl(*args: str, check: bool) -> subprocess.CompletedProcess: + return subprocess.run( + ["pactl", *args], + check=check, + capture_output=True, + text=True, encoding='utf-8', errors='replace', + stdin=subprocess.DEVNULL, + ) + + class AudioBridge: """Manages a virtual audio device for Chrome fake-mic input. - Call ``setup()`` once before launching the Meet bot and - ``teardown()`` when the session ends. ``teardown()`` is idempotent. + Call ``setup()`` once before launching the Meet bot and ``teardown()`` + when the session ends (idempotent). """ def __init__(self, name_prefix: str = "hermes_meet") -> None: @@ -42,28 +46,16 @@ class AudioBridge: self._module_ids: list[int] = [] self._torn_down = False - # ── public properties ───────────────────────────────────────────────── - - @property - def device_name(self) -> str: - if not self._device_name: + def _ready(self, value: Optional[str]) -> str: + if not value: raise RuntimeError("AudioBridge not set up yet") - return self._device_name + return value - @property - def write_target(self) -> str: - if not self._write_target: - raise RuntimeError("AudioBridge not set up yet") - return self._write_target - - # ── lifecycle ───────────────────────────────────────────────────────── + device_name = property(lambda self: self._ready(self._device_name)) + write_target = property(lambda self: self._ready(self._write_target)) def setup(self) -> dict: - """Provision the virtual audio device. - - Returns a dict describing the device. Raises RuntimeError on - unsupported platforms or when required system tools are missing. - """ + """Provision the device; raises RuntimeError on unsupported platforms or missing tools.""" system = platform.system() if system == "Linux": return self._setup_linux() @@ -74,99 +66,54 @@ class AudioBridge: raise RuntimeError(f"unsupported platform: {system}") def teardown(self) -> None: - """Release the virtual audio device. Idempotent.""" + """Release the virtual audio device. Idempotent; never raises.""" if self._torn_down: return - # Only Linux needs explicit unloading. - if self._platform == "linux" and self._module_ids: + if self._platform == "linux": # Unload in reverse order (virtual-source before null-sink). for mod_id in reversed(self._module_ids): try: - subprocess.run( - ["pactl", "unload-module", str(mod_id)], - check=False, - capture_output=True, - stdin=subprocess.DEVNULL, - ) + _pactl("unload-module", str(mod_id), check=False) except Exception: - # Best-effort teardown — never raise from here. pass self._module_ids = [] self._torn_down = True - # ── platform impls ──────────────────────────────────────────────────── + def _finish(self, tag: str, device: str, write_target: str, module_ids: list[int]) -> dict: + self._platform = tag + self._device_name = device + self._write_target = write_target + self._module_ids = module_ids + self._torn_down = False + return {"platform": tag, "device_name": device, "sample_rate": 48000, "channels": 2, + "module_ids": list(module_ids), "write_target": write_target} def _setup_linux(self) -> dict: sink_name = f"{self._name_prefix}_sink" src_name = f"{self._name_prefix}_src" try: - sink_out = subprocess.run( - [ - "pactl", - "load-module", - "module-null-sink", - f"sink_name={sink_name}", - "sink_properties=device.description=HermesMeetSink", - ], - check=True, - capture_output=True, - text=True, encoding='utf-8', errors='replace', - stdin=subprocess.DEVNULL, + sink_out = _pactl( + "load-module", "module-null-sink", f"sink_name={sink_name}", + "sink_properties=device.description=HermesMeetSink", check=True, ) except FileNotFoundError as exc: - raise RuntimeError( - "pactl not found — install PulseAudio/pipewire-pulse" - ) from exc + raise RuntimeError("pactl not found — install PulseAudio/pipewire-pulse") from exc except subprocess.CalledProcessError as exc: - raise RuntimeError( - f"pactl load-module null-sink failed: {exc.stderr or exc}" - ) from exc - + raise RuntimeError(f"pactl load-module null-sink failed: {exc.stderr or exc}") from exc sink_mod_id = self._parse_module_id(sink_out.stdout) try: - src_out = subprocess.run( - [ - "pactl", - "load-module", - "module-virtual-source", - f"source_name={src_name}", - f"master={sink_name}.monitor", - ], - check=True, - capture_output=True, - text=True, encoding='utf-8', errors='replace', - stdin=subprocess.DEVNULL, + src_out = _pactl( + "load-module", "module-virtual-source", f"source_name={src_name}", + f"master={sink_name}.monitor", check=True, ) except subprocess.CalledProcessError as exc: # Roll back the null-sink we just created so we don't leak it. - subprocess.run( - ["pactl", "unload-module", str(sink_mod_id)], - check=False, - capture_output=True, - stdin=subprocess.DEVNULL, - ) - raise RuntimeError( - f"pactl load-module virtual-source failed: {exc.stderr or exc}" - ) from exc + _pactl("unload-module", str(sink_mod_id), check=False) + raise RuntimeError(f"pactl load-module virtual-source failed: {exc.stderr or exc}") from exc - src_mod_id = self._parse_module_id(src_out.stdout) - - self._platform = "linux" - self._device_name = src_name - self._write_target = sink_name - self._module_ids = [sink_mod_id, src_mod_id] - self._torn_down = False - - return { - "platform": "linux", - "device_name": src_name, - "sample_rate": 48000, - "channels": 2, - "module_ids": list(self._module_ids), - "write_target": sink_name, - } + return self._finish("linux", src_name, sink_name, [sink_mod_id, self._parse_module_id(src_out.stdout)]) def _setup_darwin(self) -> dict: try: @@ -176,73 +123,25 @@ class AudioBridge: stderr=subprocess.STDOUT, ) except FileNotFoundError as exc: - raise RuntimeError( - "system_profiler not found (macOS-only command)" - ) from exc + raise RuntimeError("system_profiler not found (macOS-only command)") from exc except subprocess.CalledProcessError as exc: - raise RuntimeError( - f"system_profiler failed: {exc.output}" - ) from exc + raise RuntimeError(f"system_profiler failed: {exc.output}") from exc if "BlackHole" not in out: raise RuntimeError( "BlackHole virtual audio device not installed. " "Install via: brew install blackhole-2ch" ) - - self._platform = "darwin" - self._device_name = _BLACKHOLE_DEVICE - self._write_target = _BLACKHOLE_DEVICE - self._module_ids = [] - self._torn_down = False - - return { - "platform": "darwin", - "device_name": _BLACKHOLE_DEVICE, - "sample_rate": 48000, - "channels": 2, - "module_ids": [], - "write_target": _BLACKHOLE_DEVICE, - } - - # ── helpers ────────────────────────────────────────────────────────── + return self._finish("darwin", _BLACKHOLE_DEVICE, _BLACKHOLE_DEVICE, []) @staticmethod def _parse_module_id(stdout: str) -> int: - """pactl load-module prints the new module ID to stdout.""" + """pactl load-module prints the new module ID as the last token of its first line.""" text = (stdout or "").strip() if not text: raise RuntimeError("pactl load-module returned empty stdout") - # Take the last whitespace-separated token on the first non-empty line. - first = text.splitlines()[0].strip() - token = first.split()[-1] + token = text.splitlines()[0].strip().split()[-1] try: return int(token) except ValueError as exc: - raise RuntimeError( - f"could not parse pactl module id from: {stdout!r}" - ) from exc - - -def chrome_fake_audio_flags(bridge_info: dict) -> list[str]: - """Return Chrome flags for using the fake audio input. - - The PulseAudio source is selected via the ``PULSE_SOURCE`` env var, - which callers must set in Chrome's environment before launch: - - env["PULSE_SOURCE"] = bridge_info["device_name"] - - On macOS the caller must ensure the system default audio input is - set to the returned BlackHole device (we do not flip that switch). - """ - system = platform.system() - if system == "Linux": - # Chromium on Linux picks up the PulseAudio source selected via - # PULSE_SOURCE env var; the fake-ui flag skips the permission - # prompt so the bot can pick "use my mic" without user input. - return ["--use-fake-ui-for-media-stream"] - if system == "Darwin": - return ["--use-fake-ui-for-media-stream"] - if system == "Windows": - raise RuntimeError("windows not supported in v2") - raise RuntimeError(f"unsupported platform: {system}") + raise RuntimeError(f"could not parse pactl module id from: {stdout!r}") from exc diff --git a/plugins/google_meet/cli.py b/plugins/google_meet/cli.py index 45fede0381..72cf69826c 100644 --- a/plugins/google_meet/cli.py +++ b/plugins/google_meet/cli.py @@ -1,18 +1,19 @@ -"""CLI commands for the google_meet plugin. +"""CLI commands for the google_meet plugin (``hermes meet ``). -Wires ``hermes meet ``: - setup — preflight playwright, chromium, auth file, print fixes - auth — open a browser to sign into Google, save storage state - join — join a Meet URL synchronously (also callable from the agent) - status — print current bot state - transcript — print the transcript - stop — leave the current meeting + setup / install — preflight and install prerequisites + auth — open a browser to sign into Google, save storage state + join — join a Meet URL (locally or on a remote node) + status / transcript / say / stop — drive the active bot + node — remote node host management (see node/cli.py) """ from __future__ import annotations import argparse import json +import platform +import shutil +import subprocess import sys from pathlib import Path from typing import Optional @@ -21,21 +22,15 @@ from hermes_constants import get_hermes_home from plugins.google_meet import process_manager as pm from plugins.google_meet.meet_bot import _is_safe_meet_url +from plugins.google_meet.tools import resolve_node def _auth_state_path() -> Path: return Path(get_hermes_home()) / "workspace" / "meetings" / "auth.json" -# --------------------------------------------------------------------------- -# argparse wiring -# --------------------------------------------------------------------------- - def register_cli(subparser: argparse.ArgumentParser) -> None: - """Build the ``hermes meet`` argparse tree. - - Called by :func:`_register_cli_commands` at plugin load time. - """ + """Build the ``hermes meet`` argparse tree (called at plugin load time).""" subs = subparser.add_subparsers(dest="meet_command") subs.add_parser("setup", help="Preflight: playwright, chromium, auth") @@ -80,7 +75,6 @@ def register_cli(subparser: argparse.ArgumentParser) -> None: subs.add_parser("stop", help="Leave the current meeting") - # v3: remote node host management. node_p = subs.add_parser( "node", help="Manage remote meet node hosts (run/list/approve/remove/status/ping)", @@ -89,86 +83,70 @@ def register_cli(subparser: argparse.ArgumentParser) -> None: from plugins.google_meet.node.cli import register_cli as _register_node_cli _register_node_cli(node_p) except Exception as e: # pragma: no cover — defensive - # If the node module fails to import for any reason (optional dep - # missing at import time etc.), leave the subparser present but - # flag it. The argparse dispatch will surface a clear error. + # Keep the subparser present so argparse dispatch surfaces a clear error. + err = e + def _node_unavailable(args): - print(f"hermes meet node: module unavailable ({e})") + print(f"hermes meet node: module unavailable ({err})") return 1 node_p.set_defaults(func=_node_unavailable) subparser.set_defaults(func=meet_command) -# --------------------------------------------------------------------------- -# Dispatch -# --------------------------------------------------------------------------- +def _cmd_node(args: argparse.Namespace) -> int: + fn = getattr(args, "func", None) + if fn is None or fn is meet_command: + print("usage: hermes meet node {run,list,approve,remove,status,ping}") + return 2 + return fn(args) + + +_DISPATCH = { + "setup": lambda a: _cmd_setup(), + "install": lambda a: _cmd_install( + realtime=bool(getattr(a, "realtime", False)), assume_yes=bool(getattr(a, "yes", False)), + ), + "auth": lambda a: _cmd_auth(), + "join": lambda a: _cmd_join( + url=a.url, guest_name=a.guest_name, duration=a.duration, headed=a.headed, + mode=getattr(a, "mode", "transcribe"), node=getattr(a, "node", None), + ), + "status": lambda a: _print_result(pm.status()), + "transcript": lambda a: _cmd_transcript(last=a.last), + "say": lambda a: _cmd_say(text=a.text, node=getattr(a, "node", None)), + "stop": lambda a: _print_result(pm.stop(reason="hermes meet stop")), + "node": _cmd_node, +} + def meet_command(args: argparse.Namespace) -> int: sub = getattr(args, "meet_command", None) if not sub: print("usage: hermes meet {setup,auth,join,status,transcript,say,stop,node}") return 2 - if sub == "setup": - return _cmd_setup() - if sub == "install": - return _cmd_install( - realtime=bool(getattr(args, "realtime", False)), - assume_yes=bool(getattr(args, "yes", False)), - ) - if sub == "auth": - return _cmd_auth() - if sub == "join": - return _cmd_join( - url=args.url, - guest_name=args.guest_name, - duration=args.duration, - headed=args.headed, - mode=getattr(args, "mode", "transcribe"), - node=getattr(args, "node", None), - ) - if sub == "status": - return _cmd_status() - if sub == "transcript": - return _cmd_transcript(last=args.last) - if sub == "say": - return _cmd_say(text=args.text, node=getattr(args, "node", None)) - if sub == "stop": - return _cmd_stop() - if sub == "node": - # Dispatch was set by the node cli's register_cli; fall through to - # whatever its subparsers wired. - fn = getattr(args, "func", None) - if fn is None or fn is meet_command: - print("usage: hermes meet node {run,list,approve,remove,status,ping}") - return 2 - return fn(args) - print(f"unknown subcommand: {sub}") - return 2 + handler = _DISPATCH.get(sub) + if handler is None: + print(f"unknown subcommand: {sub}") + return 2 + return handler(args) -# --------------------------------------------------------------------------- -# Subcommand handlers -# --------------------------------------------------------------------------- - def _cmd_setup() -> int: - import platform as _p - print("google_meet preflight") print("---------------------") - system = _p.system() + system = platform.system() system_ok = system in {"Linux", "Darwin"} print(f" platform : {system} [{'ok' if system_ok else 'unsupported'}]") try: import playwright # noqa: F401 pw_ok = True - pw_msg = "installed" + print(" playwright : installed") except ImportError: pw_ok = False - pw_msg = "NOT installed — run: pip install playwright" - print(f" playwright : {pw_msg}") + print(" playwright : NOT installed — run: pip install playwright") chromium_ok = False chromium_msg = "unknown" @@ -176,157 +154,113 @@ def _cmd_setup() -> int: try: from playwright.sync_api import sync_playwright with sync_playwright() as p: - try: - exe = p.chromium.executable_path - if exe and Path(exe).exists(): - chromium_ok = True - chromium_msg = f"ok ({exe})" - else: - chromium_msg = ( - "not installed — run: " - "python -m playwright install chromium" - ) - except Exception as e: - chromium_msg = f"probe failed: {e}" + exe = p.chromium.executable_path + if exe and Path(exe).exists(): + chromium_ok = True + chromium_msg = f"ok ({exe})" + else: + chromium_msg = "not installed — run: python -m playwright install chromium" except Exception as e: chromium_msg = f"probe failed: {e}" print(f" chromium : {chromium_msg}") auth_path = _auth_state_path() - auth_ok = auth_path.is_file() print( " google auth : " - + (f"ok ({auth_path})" if auth_ok else "not saved — run: hermes meet auth") + + (f"ok ({auth_path})" if auth_path.is_file() else "not saved — run: hermes meet auth") ) print() all_ok = system_ok and pw_ok and chromium_ok if all_ok: - print( - "ready. Join a meeting: " - "hermes meet join https://meet.google.com/abc-defg-hij" - ) + print("ready. Join a meeting: hermes meet join https://meet.google.com/abc-defg-hij") else: print("not ready yet — fix the items above.") return 0 if all_ok else 1 def _cmd_install(*, realtime: bool, assume_yes: bool) -> int: - """Install the plugin's prerequisites. + """pip deps + Chromium; with ``--realtime`` also the platform audio bridge deps. - Always: pip install playwright + websockets, then - ``python -m playwright install chromium``. - - With ``--realtime``: also install the platform audio bridge deps. - Linux : ``sudo apt-get install -y pulseaudio-utils`` - macOS : ``brew install blackhole-2ch ffmpeg`` (+ remind the user - to select BlackHole as the default input device manually) - - Prompts before every package-manager invocation unless ``--yes``. - Refuses to run on Windows. + Prompts before every package-manager invocation unless ``--yes``. Linux/macOS only. """ - import platform as _p - import shutil as _shutil - import subprocess as _sp - - system = _p.system() + system = platform.system() if system not in {"Linux", "Darwin"}: print(f"google_meet install: {system} is not supported (linux/macos only)") return 1 - def _confirm(prompt: str) -> bool: - if assume_yes: - return True + def _install_pkgs(prompt: str, cmd: list[str], fail_msg: str) -> None: + """Confirm (unless --yes) then run a package-manager command, reporting failure.""" try: - ans = input(f"{prompt} [y/N] ").strip().lower() + ok = assume_yes or input(f"{prompt} [y/N] ").strip().lower() in {"y", "yes"} except EOFError: - return False - return ans in {"y", "yes"} + ok = False + if not ok: + print(" skipped (you can run it manually later)") + return + print(f" $ {' '.join(cmd)}") + if subprocess.run(cmd, check=False).returncode != 0: + print(fail_msg) print("google_meet install") print("-------------------") - # 1) pip deps — always safe, venv-scoped. pip_pkgs = ["playwright", "websockets"] print(f"\n[1/3] pip install: {' '.join(pip_pkgs)}") try: from hermes_cli.tools_config import _pip_install - res = _pip_install(["--upgrade", *pip_pkgs], capture_output=False) - if res.returncode != 0: + if _pip_install(["--upgrade", *pip_pkgs], capture_output=False).returncode != 0: print(" pip install failed") return 1 except Exception as e: print(f" pip install failed: {e}") return 1 - # 2) Playwright browsers — pulls chromium (~300MB first run). print("\n[2/3] python -m playwright install chromium") try: - res = _sp.run( - [sys.executable, "-m", "playwright", "install", "chromium"], - check=False, - ) + res = subprocess.run([sys.executable, "-m", "playwright", "install", "chromium"], check=False) if res.returncode != 0: print(" playwright install failed (may already be installed)") except Exception as e: print(f" playwright install failed: {e}") return 1 - # 3) Platform audio deps for realtime mode. - if realtime: + if not realtime: + print("\n[3/3] skipped (pass --realtime to install audio tooling too)") + else: print("\n[3/3] realtime audio deps") if system == "Linux": - if _shutil.which("paplay") and _shutil.which("pactl"): + if shutil.which("paplay") and shutil.which("pactl"): print(" pulseaudio-utils already installed.") else: - if not _confirm( - " install pulseaudio-utils? this runs `sudo apt-get install -y pulseaudio-utils`" - ): - print(" skipped (you can run it manually later)") - else: - cmd = ["sudo", "apt-get", "install", "-y", "pulseaudio-utils"] - print(f" $ {' '.join(cmd)}") - res = _sp.run(cmd, check=False) - if res.returncode != 0: - print(" apt install failed — install pulseaudio-utils manually") + _install_pkgs(" install pulseaudio-utils? this runs `sudo apt-get install -y pulseaudio-utils`", + ["sudo", "apt-get", "install", "-y", "pulseaudio-utils"], + " apt install failed — install pulseaudio-utils manually") elif system == "Darwin": have_bh = False try: - out = _sp.check_output(["system_profiler", "SPAudioDataType"], text=True, encoding='utf-8', errors='replace') + out = subprocess.check_output(["system_profiler", "SPAudioDataType"], text=True, encoding='utf-8', errors='replace') have_bh = "BlackHole" in out except Exception: pass - have_ffmpeg = bool(_shutil.which("ffmpeg")) - needs = [] - if not have_bh: - needs.append("blackhole-2ch") - if not have_ffmpeg: - needs.append("ffmpeg") + needs = ([] if have_bh else ["blackhole-2ch"]) + ([] if shutil.which("ffmpeg") else ["ffmpeg"]) if not needs: print(" BlackHole and ffmpeg already installed.") - elif not _shutil.which("brew"): + elif not shutil.which("brew"): print( " missing: " + ", ".join(needs) + "\n" " install Homebrew first (https://brew.sh) or install the packages manually." ) else: - if not _confirm(f" install via brew: {' '.join(needs)}?"): - print(" skipped (you can run it manually later)") - else: - cmd = ["brew", "install", *needs] - print(f" $ {' '.join(cmd)}") - res = _sp.run(cmd, check=False) - if res.returncode != 0: - print(" brew install failed — install them manually") + _install_pkgs(f" install via brew: {' '.join(needs)}?", ["brew", "install", *needs], + " brew install failed — install them manually") print( "\n NOTE: macOS does not auto-route audio. Open\n" " System Settings → Sound → Input\n" " and select 'BlackHole 2ch' before starting a realtime meeting.\n" " hermes will not switch your default input for you." ) - else: - print("\n[3/3] skipped (pass --realtime to install audio tooling too)") print("\ndone. verify with: hermes meet setup") return 0 @@ -367,6 +301,29 @@ def _cmd_auth() -> int: return 0 +def _print_result(res: dict) -> int: + print(json.dumps(res, indent=2)) + return 0 if res.get("ok") else 1 + + +def _remote(node: str, op: str, call) -> int: + """Run *call(client)* against the registered node *node* and print the result.""" + try: + client, name = resolve_node(node) + except ImportError as e: + print(f"node module unavailable: {e}") + return 1 + if client is None: + print(f"no registered node matches {node!r}") + return 1 + try: + res = call(client) + except Exception as e: + print(f"remote {op} failed: {e}") + return 1 + return _print_result({"node": name, **res}) + + def _cmd_join( url: str, *, @@ -380,41 +337,14 @@ def _cmd_join( print(f"refusing: not a meet.google.com URL: {url}") return 2 if node: - # Remote: go through NodeClient. - try: - from plugins.google_meet.node.registry import NodeRegistry - from plugins.google_meet.node.client import NodeClient - except ImportError as e: - print(f"node module unavailable: {e}") - return 1 - reg = NodeRegistry() - entry = reg.resolve(node if node != "auto" else None) - if entry is None: - print(f"no registered node matches {node!r}") - return 1 - client = NodeClient(url=entry["url"], token=entry["token"]) - try: - res = client.start_bot( - url=url, guest_name=guest_name, duration=duration, - headed=headed, mode=mode, - ) - except Exception as e: - print(f"remote start_bot failed: {e}") - return 1 - print(json.dumps({"node": entry.get("name"), **res}, indent=2)) - return 0 if res.get("ok") else 1 - + return _remote(node, "start_bot", lambda c: c.start_bot( + url=url, guest_name=guest_name, duration=duration, headed=headed, mode=mode, + )) auth = _auth_state_path() - res = pm.start( - url=url, - headed=headed, - guest_name=guest_name, - duration=duration, - auth_state=str(auth) if auth.is_file() else None, - mode=mode, - ) - print(json.dumps(res, indent=2)) - return 0 if res.get("ok") else 1 + return _print_result(pm.start( + url=url, headed=headed, guest_name=guest_name, duration=duration, + auth_state=str(auth) if auth.is_file() else None, mode=mode, + )) def _cmd_say(text: str, node: Optional[str] = None) -> int: @@ -422,55 +352,20 @@ def _cmd_say(text: str, node: Optional[str] = None) -> int: print("refusing: empty text") return 2 if node: - try: - from plugins.google_meet.node.registry import NodeRegistry - from plugins.google_meet.node.client import NodeClient - except ImportError as e: - print(f"node module unavailable: {e}") - return 1 - reg = NodeRegistry() - entry = reg.resolve(node if node != "auto" else None) - if entry is None: - print(f"no registered node matches {node!r}") - return 1 - client = NodeClient(url=entry["url"], token=entry["token"]) - try: - res = client.say(text) - except Exception as e: - print(f"remote say failed: {e}") - return 1 - print(json.dumps({"node": entry.get("name"), **res}, indent=2)) - return 0 if res.get("ok") else 1 - - res = pm.enqueue_say(text) - print(json.dumps(res, indent=2)) - return 0 if res.get("ok") else 1 - - -def _cmd_status() -> int: - res = pm.status() - print(json.dumps(res, indent=2)) - return 0 if res.get("ok") else 1 + return _remote(node, "say", lambda c: c.say(text)) + return _print_result(pm.enqueue_say(text)) def _cmd_transcript(last: Optional[int]) -> int: res = pm.transcript(last=last) if not res.get("ok"): - print(json.dumps(res, indent=2)) - return 1 + return _print_result(res) for ln in res.get("lines", []): print(ln) return 0 -def _cmd_stop() -> int: - res = pm.stop(reason="hermes meet stop") - print(json.dumps(res, indent=2)) - return 0 if res.get("ok") else 1 - - if __name__ == "__main__": # pragma: no cover parser = argparse.ArgumentParser(prog="hermes meet") register_cli(parser) - ns = parser.parse_args() - sys.exit(meet_command(ns)) + sys.exit(meet_command(parser.parse_args())) diff --git a/plugins/google_meet/meet_bot.py b/plugins/google_meet/meet_bot.py index 20745eb3de..af89a00a64 100644 --- a/plugins/google_meet/meet_bot.py +++ b/plugins/google_meet/meet_bot.py @@ -1,45 +1,35 @@ """Headless Google Meet bot — Playwright + live-caption scraping. -Runs as a standalone subprocess spawned by ``process_manager.py``. Reads config -from env vars, writes status + transcript to files under -``$HERMES_HOME/workspace/meetings//``. The main hermes process -reads those files via the ``meet_*`` tools — no IPC beyond filesystem. +Standalone subprocess spawned by ``process_manager.py``. Config comes from env +vars; status + transcript are written under ``$HERMES_MEET_OUT_DIR`` and read +by the ``meet_*`` tools — no IPC beyond the filesystem. -The scraping strategy mirrors OpenUtter (sumansid/openutter): we don't parse -WebRTC audio, we enable Google Meet's built-in live captions and observe the -captions container in the DOM via a MutationObserver. This is lossy and -English-biased but it is: +We don't parse WebRTC audio: we enable Meet's built-in live captions and watch +the caption container via a MutationObserver. Lossy and English-biased, but +deterministic (no STT billing) and stable across Meet UI rewrites thanks to +the container's ARIA role. Only ``https://meet.google.com/`` URLs are accepted. -* deterministic (no API keys, no STT billing), -* works behind Meet's normal login / admission, -* survives Meet UI rewrites fairly well because the caption container has a - stable ARIA role. - -Run standalone for debugging:: - - HERMES_MEET_URL=https://meet.google.com/abc-defg-hij \\ - HERMES_MEET_OUT_DIR=/tmp/meet-debug \\ - HERMES_MEET_HEADED=1 \\ - python -m plugins.google_meet.meet_bot - -No meet.google.com URL → exits non-zero. Any URL that doesn't start with -``https://meet.google.com/`` is rejected (explicit-by-design). +Debug run: ``HERMES_MEET_URL=... HERMES_MEET_OUT_DIR=/tmp/meet-debug +HERMES_MEET_HEADED=1 python -m plugins.google_meet.meet_bot`` """ from __future__ import annotations -import json import os import re +import shutil import signal +import subprocess import sys import threading import time +from dataclasses import dataclass from pathlib import Path from typing import Optional -# Match ``https://meet.google.com/abc-defg-hij`` or ``.../lookup/...`` — the -# short three-segment code or a lookup URL. Anything else is rejected. +from plugins.google_meet._jsonfile import write_json_atomic + +# Short three-segment code, a lookup URL, or /new. Anything else is rejected. MEET_URL_RE = re.compile( r"^https://meet\.google\.com/(" r"[a-z0-9]{3,}-[a-z0-9]{3,}-[a-z0-9]{3,}" @@ -48,124 +38,86 @@ MEET_URL_RE = re.compile( r")(?:[/?#].*)?$" ) - # Filenames the bot reads/writes in ``HERMES_MEET_OUT_DIR``. SAY_QUEUE_FILENAME = "say_queue.jsonl" SAY_PCM_FILENAME = "speaker.pcm" +_FFMPEG_MISSING = "ffmpeg not found — install via `brew install ffmpeg` for realtime on macOS" + def _is_safe_meet_url(url: str) -> bool: - """Return True if *url* is a Google Meet URL we're willing to navigate to.""" - if not isinstance(url, str): - return False - return bool(MEET_URL_RE.match(url.strip())) + """True if *url* is a Google Meet URL we're willing to navigate to.""" + return isinstance(url, str) and bool(MEET_URL_RE.match(url.strip())) def _meeting_id_from_url(url: str) -> str: - """Extract the 3-segment meeting code from a Meet URL. - - For ``https://meet.google.com/abc-defg-hij`` → ``abc-defg-hij``. - For ``.../lookup/`` or ``/new`` we fall back to a timestamped id — the - bot won't know the real code until after redirect, and callers pass this - through to filename anyway. - """ - m = re.search( - r"meet\.google\.com/([a-z0-9]{3,}-[a-z0-9]{3,}-[a-z0-9]{3,})", - url or "", - ) - if m: - return m.group(1) - return f"meet-{int(time.time())}" + """3-segment meeting code, or a timestamped id for ``/lookup/...`` and ``/new``.""" + m = re.search(r"meet\.google\.com/([a-z0-9]{3,}-[a-z0-9]{3,}-[a-z0-9]{3,})", url or "") + return m.group(1) if m else f"meet-{int(time.time())}" -# --------------------------------------------------------------------------- -# Status + transcript file writers -# --------------------------------------------------------------------------- +def _quiet(fn, *args, **kwargs): + """Call *fn*, swallowing any exception (best-effort teardown steps).""" + try: + return fn(*args, **kwargs) + except Exception: + return None + + +# status.json keys in file order → _BotState attribute + initial value. +_STATUS_FIELDS = ( + ("meetingId", "meeting_id", None), ("url", "url", None), + ("inCall", "in_call", False), ("captioning", "captioning", False), + ("captionsEnabledAttempted", "captions_enabled_attempted", False), + ("lobbyWaiting", "lobby_waiting", False), + ("joinAttemptedAt", "join_attempted_at", None), ("joinedAt", "joined_at", None), + ("lastCaptionAt", "last_caption_at", None), ("transcriptLines", "transcript_lines", 0), + ("transcriptPath", "transcript_path", None), + ("error", "error", None), ("exited", "exited", False), ("pid", None, None), + # v2 realtime telemetry. + ("realtime", "realtime", False), ("realtimeReady", "realtime_ready", False), + ("realtimeDevice", "realtime_device", None), ("audioBytesOut", "audio_bytes_out", 0), + ("lastAudioOutAt", "last_audio_out_at", None), ("lastBargeInAt", "last_barge_in_at", None), + ("leaveReason", "leave_reason", None), +) + class _BotState: """Single-process mutable state, flushed to ``status.json`` on each change.""" def __init__(self, out_dir: Path, meeting_id: str, url: str): + for _, attr, default in _STATUS_FIELDS: + if attr: + setattr(self, attr, default) self.out_dir = out_dir self.meeting_id = meeting_id self.url = url - self.in_call = False - self.captioning = False - self.captions_enabled_attempted = False - self.lobby_waiting = False - self.join_attempted_at: Optional[float] = None - self.joined_at: Optional[float] = None - self.last_caption_at: Optional[float] = None - self.transcript_lines = 0 - self.error: Optional[str] = None - self.exited = False - # v2 realtime fields. - self.realtime = False - self.realtime_ready = False - self.realtime_device: Optional[str] = None - self.audio_bytes_out: int = 0 - self.last_audio_out_at: Optional[float] = None - self.last_barge_in_at: Optional[float] = None - self.leave_reason: Optional[str] = None - # Scraped captions, in order, deduped. Each entry is a dict of - # {"ts": , "speaker": str, "text": str}. - self._seen: set = set() + self._seen: set = set() # "speaker|text" keys already written out_dir.mkdir(parents=True, exist_ok=True) self.transcript_path = out_dir / "transcript.txt" self.status_path = out_dir / "status.json" self._flush() - # -------- transcript ------------------------------------------------ - def record_caption(self, speaker: str, text: str) -> None: - """Append a caption line if we haven't seen this exact (speaker, text).""" + """Append a caption line unless this exact (speaker, text) was already seen.""" speaker = (speaker or "").strip() or "Unknown" text = (text or "").strip() - if not text: - return key = f"{speaker}|{text}" - if key in self._seen: + if not text or key in self._seen: return self._seen.add(key) self.transcript_lines += 1 self.last_caption_at = time.time() ts = time.strftime("%H:%M:%S", time.localtime(self.last_caption_at)) - line = f"[{ts}] {speaker}: {text}\n" - # Atomic-ish append — good enough for a single-writer. with self.transcript_path.open("a", encoding="utf-8") as f: - f.write(line) + f.write(f"[{ts}] {speaker}: {text}\n") self._flush() - # -------- status file ---------------------------------------------- - def _flush(self) -> None: - data = { - "meetingId": self.meeting_id, - "url": self.url, - "inCall": self.in_call, - "captioning": self.captioning, - "captionsEnabledAttempted": self.captions_enabled_attempted, - "lobbyWaiting": self.lobby_waiting, - "joinAttemptedAt": self.join_attempted_at, - "joinedAt": self.joined_at, - "lastCaptionAt": self.last_caption_at, - "transcriptLines": self.transcript_lines, - "transcriptPath": str(self.transcript_path), - "error": self.error, - "exited": self.exited, - "pid": os.getpid(), - # v2 realtime telemetry. - "realtime": self.realtime, - "realtimeReady": self.realtime_ready, - "realtimeDevice": self.realtime_device, - "audioBytesOut": self.audio_bytes_out, - "lastAudioOutAt": self.last_audio_out_at, - "lastBargeInAt": self.last_barge_in_at, - "leaveReason": self.leave_reason, - } - tmp = self.status_path.with_suffix(".json.tmp") - tmp.write_text(json.dumps(data, indent=2), encoding="utf-8") - tmp.replace(self.status_path) + data = {key: getattr(self, attr) if attr else None for key, attr, _ in _STATUS_FIELDS} + data["transcriptPath"] = str(self.transcript_path) + data["pid"] = os.getpid() # overrides keep the table's key order + write_json_atomic(self.status_path, data) def set(self, **kwargs) -> None: for k, v in kwargs.items(): @@ -173,14 +125,8 @@ class _BotState: self._flush() -# --------------------------------------------------------------------------- -# Playwright bot entry point -# --------------------------------------------------------------------------- - -# JavaScript injected into the Meet tab to observe captions. Captures -# {speaker, text} tuples via a MutationObserver on the caption container, -# and exposes ``window.__hermesMeetDrain()`` to pull new entries. This -# mirrors the OpenUtter caption scraping approach. +# JS injected into the Meet tab: MutationObserver on the caption container +# collects {speaker, text}; ``window.__hermesMeetDrain()`` pulls new entries. _CAPTION_OBSERVER_JS = r""" (() => { if (window.__hermesMeetInstalled) return; @@ -201,36 +147,31 @@ _CAPTION_OBSERVER_JS = r""" } function scan(root) { - // Meet captions render as a list of rows; each row contains a speaker - // label and a text block. Selectors vary across Meet rewrites; we try - // a few shapes and fall back to raw text. + // Meet captions render as rows of speaker label + text block. Selectors + // vary across Meet rewrites; try a few shapes and fall back to raw text. const rows = root.querySelectorAll('div[jsname="dsyhDe"], div.CNusmb, div.TBMuR'); if (rows.length) { rows.forEach((row) => { const spkEl = row.querySelector('div.KcIKyf, div.zs7s8d, span[jsname="YSxPC"]'); const txtEl = row.querySelector('div.bh44bd, span[jsname="tgaKEf"], div.iTTPOb'); - const speaker = spkEl ? spkEl.innerText : ''; - const text = txtEl ? txtEl.innerText : row.innerText; - pushEntry(speaker, text); + pushEntry(spkEl ? spkEl.innerText : '', txtEl ? txtEl.innerText : row.innerText); }); return; } // Fallback: treat the whole region's innerText as one anonymous line. - const text = (root.innerText || '').split('\n').filter(Boolean).pop(); - pushEntry('', text); + pushEntry('', (root.innerText || '').split('\n').filter(Boolean).pop()); } function attach() { const el = document.querySelector(captionSelector); if (!el) return false; - const obs = new MutationObserver(() => scan(el)); - obs.observe(el, { childList: true, subtree: true, characterData: true }); + new MutationObserver(() => scan(el)).observe(el, { childList: true, subtree: true, characterData: true }); scan(el); return true; } - // Try now and retry on interval — the caption region only appears after - // captions are enabled and someone speaks. + // Retry on interval — the caption region only appears after captions are + // enabled and someone speaks. if (!attach()) { const iv = setInterval(() => { if (attach()) clearInterval(iv); }, 1500); } @@ -243,242 +184,313 @@ _CAPTION_OBSERVER_JS = r""" })(); """ +# Best-effort caption toggle: Meet binds it to the ``c`` key; click targeting +# is too brittle to rely on. +_ENABLE_CAPTIONS_JS = ( + "(() => { document.body.dispatchEvent(new KeyboardEvent('keydown', " + "{ key: 'c', code: 'KeyC', keyCode: 67, which: 67, bubbles: true })); return true; })();" +) -def _enable_captions_js() -> str: - """Return a small JS snippet that tries to click the 'Turn on captions' button. +_LEAVE_CALL_JS = ( + "() => { const b = document.querySelector('button[aria-label*=\"eave call\"]');" + " if (b) b.click(); }" +) - Best-effort — Meet's caption toggle is keyboard-accessible via ``c``. We - dispatch that keystroke as a cheap fallback. Real click targeting is too - brittle to rely on. - """ - return r""" +# True once we're clearly past the lobby: leave button, caption region +# (only once our observer is installed) or participant list visible. +_ADMISSION_PROBE_JS = r""" (() => { - const ev = new KeyboardEvent('keydown', { - key: 'c', code: 'KeyC', keyCode: 67, which: 67, bubbles: true, - }); - document.body.dispatchEvent(ev); - return true; + if (document.querySelector('button[aria-label*="eave call" i]')) return true; + if (window.__hermesMeetInstalled && document.querySelector( + '[role="region"][aria-label*="aption" i], div[jsname="YSxPC"], div[jsname="tgaKEf"]')) return true; + return !!document.querySelector('[aria-label*="articipants" i]'); + })(); + """ + +# English only — what Meet shows when the host denies or removes a guest. +_DENIED_PROBE_JS = r""" + (() => { + const text = document.body ? document.body.innerText || '' : ''; + return /You can't join this video call|You were removed from the meeting|No one responded to your request to join/i.test(text); })(); """ -def _start_realtime_speaker( - *, - rt: dict, - out_dir: Path, - bridge_info: dict, - api_key: str, - model: str, - voice: str, - instructions: str, - stop_flag: dict, - state: "_BotState", -) -> None: - """Wire up the OpenAI Realtime session + speaker thread + PCM pump. - - The speaker thread reads text lines from ``say_queue.jsonl``, sends each - to OpenAI Realtime, and writes PCM audio into ``speaker.pcm``. A - separate *pump* thread forwards that PCM into the OS audio sink so - Chrome's fake mic picks it up. On Linux we pipe to ``paplay`` against - the null-sink; on macOS the caller is expected to have the BlackHole - device selected as default input. - """ +def _probe(page, js: str) -> bool: + """Evaluate a boolean JS probe; conservative — False on any error.""" try: - from plugins.google_meet.realtime.openai_client import ( - RealtimeSession, - RealtimeSpeaker, + return bool(page.evaluate(js)) + except Exception: + return False + + +def _visible(locator): + """``locator.first`` if it exists and is visible, else None (swallows Playwright errors).""" + try: + first = locator.first + return first if first.count() and first.is_visible() else None + except Exception: + return None + + +def _start_pcm_pump(rt: dict, bridge_info: dict, pcm_path: Path, state: "_BotState") -> None: + """Stream the growing ``speaker.pcm`` (24kHz s16le mono) into the OS device Chrome's fake mic reads.""" + bridge_info = bridge_info or {} + platform_tag = bridge_info.get("platform") + if platform_tag == "linux": + sink = bridge_info.get("write_target") or "hermes_meet_sink" + cmd = ["paplay", "--raw", "--rate=24000", "--format=s16le", "--channels=1", + f"--device={sink}", str(pcm_path)] + missing = "paplay not found — install pulseaudio-utils for realtime on Linux" + elif platform_tag == "darwin": + # The user must have BlackHole selected as default input for Chrome + # to pick it up; ffmpeg targets the device by audiotoolbox index. + if not shutil.which("ffmpeg"): + state.set(error=_FFMPEG_MISSING) + return + device_name = bridge_info.get("write_target") or "BlackHole 2ch" + cmd = ["ffmpeg", "-nostdin", "-hide_banner", "-loglevel", "error", "-re", + "-f", "s16le", "-ar", "24000", "-ac", "1", "-i", str(pcm_path), + "-f", "audiotoolbox", "-audio_device_index", _mac_audio_device_index(device_name), "-"] + missing = _FFMPEG_MISSING + else: + return + try: + rt["pcm_pump"] = subprocess.Popen( + cmd, stdin=subprocess.DEVNULL, stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL, ) + except FileNotFoundError: + state.set(error=missing) + except Exception as e: + if platform_tag != "darwin": + raise + state.set(error=f"macOS pcm pump failed to start: {e}") + + +def _start_realtime_speaker(rt: dict, cfg: "_BotConfig", stop_flag: dict, state: "_BotState") -> None: + """Wire up the OpenAI Realtime session, the say-queue speaker thread and the PCM pump.""" + try: + from plugins.google_meet.realtime.openai_client import RealtimeSession, RealtimeSpeaker except Exception as e: state.set(error=f"realtime import failed: {e}") return - pcm_path = out_dir / SAY_PCM_FILENAME - queue_path = out_dir / SAY_QUEUE_FILENAME - processed_path = out_dir / "say_processed.jsonl" - # Reset the sink file so we start clean each session. - pcm_path.write_bytes(b"") - # Make sure the queue exists so the speaker poller doesn't error on - # first iteration. - queue_path.touch() + pcm_path = cfg.out_dir / SAY_PCM_FILENAME + queue_path = cfg.out_dir / SAY_QUEUE_FILENAME + pcm_path.write_bytes(b"") # start each session with a clean sink file + queue_path.touch() # so the speaker poller doesn't error on first iteration try: session = RealtimeSession( - api_key=api_key, - model=model, - voice=voice, - instructions=instructions, - audio_sink_path=pcm_path, - sample_rate=24000, + api_key=cfg.realtime_api_key, model=cfg.realtime_model, voice=cfg.realtime_voice, + instructions=cfg.realtime_instructions, audio_sink_path=pcm_path, sample_rate=24000, ) session.connect() except Exception as e: state.set(error=f"realtime connect failed: {e}") return - rt["session"] = session - def _stop_fn(): - return stop_flag.get("stop", False) - - rt["speaker_stop"] = lambda: stop_flag.__setitem__("stop", stop_flag.get("stop", False)) - speaker = RealtimeSpeaker( - session=session, - queue_path=queue_path, - processed_path=processed_path, + session=session, queue_path=queue_path, processed_path=cfg.out_dir / "say_processed.jsonl", ) def _speaker_loop(): try: - speaker.run_until_stopped(_stop_fn) + speaker.run_until_stopped(lambda: stop_flag.get("stop", False)) except Exception as e: state.set(error=f"realtime speaker crashed: {e}") - t_speaker = threading.Thread(target=_speaker_loop, name="meet-speaker", daemon=True) - t_speaker.start() - rt["speaker_thread"] = t_speaker - - # PCM pump: feeds speaker.pcm (24kHz s16le mono) into the OS audio - # device that Chrome's fake mic reads from. Different tools per - # platform, but the contract is the same — block-read the growing - # PCM file and stream it to the device in near-real-time. - platform_tag = (bridge_info or {}).get("platform") - if platform_tag == "linux": - import subprocess as _sp - - sink = (bridge_info or {}).get("write_target") or "hermes_meet_sink" - try: - proc = _sp.Popen( - [ - "paplay", - "--raw", - "--rate=24000", - "--format=s16le", - "--channels=1", - f"--device={sink}", - str(pcm_path), - ], - stdin=_sp.DEVNULL, - stdout=_sp.DEVNULL, - stderr=_sp.DEVNULL, - ) - rt["pcm_pump"] = proc - except FileNotFoundError: - state.set(error="paplay not found — install pulseaudio-utils for realtime on Linux") - elif platform_tag == "darwin": - # macOS: use ffmpeg to tail-read speaker.pcm and write it to the - # BlackHole output device. The user must have BlackHole selected - # as the default input in System Settings → Sound for Chrome to - # pick it up. We prefer ffmpeg because it's scriptable and can - # target AVFoundation devices by name; fall back to afplay-ing - # the file in a tight loop if ffmpeg is absent. - import shutil as _shutil - import subprocess as _sp - - device_name = (bridge_info or {}).get("write_target") or "BlackHole 2ch" - if _shutil.which("ffmpeg"): - try: - # -re: read input at native frame rate. - # -f avfoundation -i: speaker path as raw PCM. - # -f s16le -ar 24000 -ac 1 -i : interpret the file. - # -f audiotoolbox -audio_device_index: write to BlackHole. - # Simpler: output as raw via coreaudio using "-f audiotoolbox". - # ffmpeg's audiotoolbox output picks the current default - # output device, which isn't what we want. Instead we use - # -f avfoundation with the named device as OUTPUT via - # -vn and the device name. - proc = _sp.Popen( - [ - "ffmpeg", - "-nostdin", "-hide_banner", "-loglevel", "error", - "-re", - "-f", "s16le", "-ar", "24000", "-ac", "1", - "-i", str(pcm_path), - "-f", "audiotoolbox", - "-audio_device_index", _mac_audio_device_index(device_name), - "-", - ], - stdin=_sp.DEVNULL, - stdout=_sp.DEVNULL, - stderr=_sp.DEVNULL, - ) - rt["pcm_pump"] = proc - except FileNotFoundError: - state.set(error="ffmpeg not found — install via `brew install ffmpeg` for realtime on macOS") - except Exception as e: - state.set(error=f"macOS pcm pump failed to start: {e}") - else: - state.set(error="ffmpeg not found — install via `brew install ffmpeg` for realtime on macOS") + rt["speaker_thread"] = threading.Thread(target=_speaker_loop, name="meet-speaker", daemon=True) + rt["speaker_thread"].start() + _start_pcm_pump(rt, rt["bridge_info"], pcm_path, state) + state.set(realtime_ready=True) def _mac_audio_device_index(device_name: str) -> str: - """Return the ffmpeg ``-audio_device_index`` for *device_name*, as a string. + """ffmpeg ``-audio_device_index`` for *device_name* (case-insensitive), ``"0"`` if not found. - Probes ``ffmpeg -f avfoundation -list_devices true -i ''`` (which prints - the device table on stderr) and matches *device_name* case-insensitively. - Defaults to ``"0"`` if the device can't be found — caller will get a - misrouted stream but not a crash, and the error will be obvious. + ffmpeg prints the avfoundation device table on stderr as ``[N] Name``. """ - import subprocess as _sp - try: - out = _sp.run( + out = subprocess.run( ["ffmpeg", "-f", "avfoundation", "-list_devices", "true", "-i", ""], - capture_output=True, - text=True, encoding='utf-8', errors='replace', - timeout=10, + capture_output=True, text=True, encoding='utf-8', errors='replace', timeout=10, ) except Exception: return "0" - # ffmpeg prints the table on stderr. Lines look like: - # [AVFoundation indev @ 0x...] [0] BlackHole 2ch - import re as _re - needle = device_name.strip().lower() for line in (out.stderr or "").splitlines(): - m = _re.search(r"\[(\d+)\]\s+(.+)$", line) - if not m: - continue - if m.group(2).strip().lower() == needle: + m = re.search(r"\[(\d+)\]\s+(.+)$", line) + if m and m.group(2).strip().lower() == needle: return m.group(1) return "0" -def run_bot() -> int: # noqa: C901 — orchestration, explicit branches - url = os.environ.get("HERMES_MEET_URL", "").strip() - out_dir_env = os.environ.get("HERMES_MEET_OUT_DIR", "").strip() - headed = os.environ.get("HERMES_MEET_HEADED", "").lower() in {"1", "true", "yes"} - auth_state = os.environ.get("HERMES_MEET_AUTH_STATE", "").strip() - guest_name = os.environ.get("HERMES_MEET_GUEST_NAME", "Hermes Agent") - duration_s = _parse_duration(os.environ.get("HERMES_MEET_DURATION", "")) - # v2: optional realtime mode. Enabled when HERMES_MEET_MODE=realtime. - mode = os.environ.get("HERMES_MEET_MODE", "transcribe").strip().lower() - realtime_model = os.environ.get("HERMES_MEET_REALTIME_MODEL", "gpt-realtime") - realtime_voice = os.environ.get("HERMES_MEET_REALTIME_VOICE", "alloy") - realtime_instructions = os.environ.get("HERMES_MEET_REALTIME_INSTRUCTIONS", "") - # HERMES_MEET_REALTIME_KEY is set explicitly by process_manager.start(), - # which resolves it through the parent's profile secret scope at spawn - # time. The bare OPENAI_API_KEY fallback only serves standalone - # `python -m plugins.google_meet.meet_bot` runs outside the gateway. - realtime_api_key = os.environ.get("HERMES_MEET_REALTIME_KEY") or os.environ.get("OPENAI_API_KEY", "") +def _setup_realtime(rt: dict, api_key: str, state: _BotState) -> None: + """Provision the virtual audio bridge; on any failure fall back to transcribe mode.""" + if not api_key: + state.set(error="realtime mode requested but no API key in HERMES_MEET_REALTIME_KEY/OPENAI_API_KEY — falling back to transcribe") + rt["enabled"] = False + return + try: + from plugins.google_meet.audio_bridge import AudioBridge + bridge = AudioBridge() + rt["bridge_info"] = bridge.setup() + rt["bridge"] = bridge + state.set(realtime=True, realtime_device=rt["bridge_info"].get("device_name")) + except Exception as e: + state.set(error=f"audio bridge setup failed: {e} — falling back to transcribe") + rt["enabled"] = False - if not url or not _is_safe_meet_url(url): + +def _teardown_realtime(rt: dict) -> None: + if rt.get("pcm_pump"): + _quiet(rt["pcm_pump"].terminate) + _quiet(rt["pcm_pump"].wait, timeout=3) + if rt["speaker_thread"] is not None: + _quiet(rt["speaker_thread"].join, timeout=5.0) + if rt["session"]: + _quiet(rt["session"].close) + if rt["bridge"]: + _quiet(rt["bridge"].teardown) + + +@dataclass +class _BotConfig: + """Everything the bot reads from ``HERMES_MEET_*`` env vars.""" + + url: str + out_dir: Optional[Path] + headed: bool + auth_state: str + guest_name: str + duration_s: Optional[float] + realtime: bool + realtime_api_key: str + realtime_model: str + realtime_voice: str + realtime_instructions: str + lobby_timeout: float + + +def _config_from_env() -> _BotConfig: + env = os.environ.get + out_raw = env("HERMES_MEET_OUT_DIR", "").strip() + return _BotConfig( + url=env("HERMES_MEET_URL", "").strip(), + out_dir=Path(out_raw) if out_raw else None, + headed=env("HERMES_MEET_HEADED", "").lower() in {"1", "true", "yes"}, + auth_state=env("HERMES_MEET_AUTH_STATE", "").strip(), + guest_name=env("HERMES_MEET_GUEST_NAME", "Hermes Agent"), + duration_s=_parse_duration(env("HERMES_MEET_DURATION", "")), + realtime=env("HERMES_MEET_MODE", "transcribe").strip().lower() == "realtime", + # HERMES_MEET_REALTIME_KEY is resolved by process_manager.start() through the + # parent's profile secret scope; the OPENAI_API_KEY fallback only serves + # standalone `python -m plugins.google_meet.meet_bot` runs. + realtime_api_key=env("HERMES_MEET_REALTIME_KEY") or env("OPENAI_API_KEY", ""), + realtime_model=env("HERMES_MEET_REALTIME_MODEL", "gpt-realtime"), + realtime_voice=env("HERMES_MEET_REALTIME_VOICE", "alloy"), + realtime_instructions=env("HERMES_MEET_REALTIME_INSTRUCTIONS", ""), + lobby_timeout=float(env("HERMES_MEET_LOBBY_TIMEOUT", "300")), + ) + + +def _join(page, cfg: _BotConfig, state: _BotState) -> None: + """Fill the guest-name field (guest mode) and click 'Join now' / 'Ask to join'. + + 'Ask to join' means we're in the lobby → ``lobby_waiting``. + """ + name_box = _visible(page.locator('input[aria-label*="name" i]')) + if name_box is not None: + _quiet(name_box.fill, cfg.guest_name, timeout=2_000) + for label in ("Join now", "Ask to join"): + btn = _visible(page.get_by_role("button", name=label, exact=False)) + if btn is None: + continue + try: + btn.click(timeout=3_000) + if label == "Ask to join": + state.set(lobby_waiting=True) + break + except Exception: + continue + + +def _drain_loop(page, cfg: _BotConfig, state: _BotState, rt: dict, stop_flag: dict) -> None: + """Admission + caption drain loop; runs until SIGTERM, duration expiry, lobby timeout/denial or page loss. + + Sets ``state.leave_reason`` for every exit but SIGTERM. Also triggers + barge-in and mirrors realtime counters into status.json. + """ + deadline = (time.time() + cfg.duration_s) if cfg.duration_s else None + lobby_deadline = time.time() + cfg.lobby_timeout + last_admission_check = 0.0 + while not stop_flag["stop"]: + now = time.time() + if deadline and now > deadline: + state.set(leave_reason="duration_expired") + return + + if not state.in_call and (now - last_admission_check) > 3.0: + last_admission_check = now + if _probe(page, _ADMISSION_PROBE_JS): + state.set(in_call=True, lobby_waiting=False, joined_at=now) + elif now > lobby_deadline: + waited = int(lobby_deadline - state.join_attempted_at) if state.join_attempted_at else 0 + state.set( + error=f"lobby timeout — host never admitted the bot within {waited}s", + leave_reason="lobby_timeout", + ) + return + elif _probe(page, _DENIED_PROBE_JS): + state.set(error="host denied admission", leave_reason="denied") + return + + try: + queued = page.evaluate("window.__hermesMeetDrain && window.__hermesMeetDrain()") + for entry in queued if isinstance(queued, list) else (): + if not isinstance(entry, dict): + continue + speaker = str(entry.get("speaker", "")) + state.record_caption(speaker=speaker, text=str(entry.get("text", ""))) + # Barge-in: a real human spoke while we may be generating + # audio — cancel the in-flight response. + if (rt["session"] is not None + and _looks_like_human_speaker(speaker, cfg.guest_name) + and _quiet(rt["session"].cancel_response)): + state.set(last_barge_in_at=now) + except Exception: + # Meet reloaded or we got booted — exit rather than spin. + if page.is_closed(): + state.set(leave_reason="page_closed") + return + + if rt["session"] is not None: + state.set( + audio_bytes_out=rt["session"].audio_bytes_out, + last_audio_out_at=rt["session"].last_audio_out_at, + ) + + time.sleep(1.0) + + +def run_bot() -> int: + cfg = _config_from_env() + if not _is_safe_meet_url(cfg.url): sys.stderr.write( "google_meet bot: refusing to launch — HERMES_MEET_URL must be a " - "meet.google.com URL. got: %r\n" % url + "meet.google.com URL. got: %r\n" % cfg.url ) return 2 - if not out_dir_env: + if cfg.out_dir is None: sys.stderr.write("google_meet bot: HERMES_MEET_OUT_DIR is required\n") return 2 - out_dir = Path(out_dir_env) - meeting_id = _meeting_id_from_url(url) - state = _BotState(out_dir=out_dir, meeting_id=meeting_id, url=url) + state = _BotState(out_dir=cfg.out_dir, meeting_id=_meeting_id_from_url(cfg.url), url=cfg.url) - # SIGTERM → exit cleanly so the parent ``meet_leave`` gets a finalized - # transcript. We set a flag instead of raising so the Playwright context - # teardown runs in the finally block below. + # SIGTERM sets a flag (not an exception) so the Playwright teardown below + # still runs and ``meet_leave`` gets a finalized transcript. stop_flag = {"stop": False} def _on_signal(_sig, _frame): @@ -487,32 +499,11 @@ def run_bot() -> int: # noqa: C901 — orchestration, explicit branches signal.signal(signal.SIGTERM, _on_signal) signal.signal(signal.SIGINT, _on_signal) - # v2 realtime: provision virtual audio device + start speaker thread. - # We track these in a dict so the finally block can tear them down - # regardless of how we exit. If anything in the realtime setup fails we - # fall back to transcribe mode with a status flag. - rt = { - "enabled": mode == "realtime", - "bridge": None, # AudioBridge | None - "bridge_info": None, # dict | None - "session": None, # RealtimeSession | None - "speaker_thread": None, # threading.Thread | None - "speaker_stop": None, # callable | None - } + # Realtime resources tracked in one dict so teardown works however we exit. + rt = {"enabled": cfg.realtime, "bridge": None, "bridge_info": None, + "session": None, "speaker_thread": None} if rt["enabled"]: - if not realtime_api_key: - state.set(error="realtime mode requested but no API key in HERMES_MEET_REALTIME_KEY/OPENAI_API_KEY — falling back to transcribe") - rt["enabled"] = False - else: - try: - from plugins.google_meet.audio_bridge import AudioBridge - bridge = AudioBridge() - rt["bridge_info"] = bridge.setup() - rt["bridge"] = bridge - state.set(realtime=True, realtime_device=rt["bridge_info"].get("device_name")) - except Exception as e: - state.set(error=f"audio bridge setup failed: {e} — falling back to transcribe") - rt["enabled"] = False + _setup_realtime(rt, cfg.realtime_api_key, state) try: from playwright.sync_api import sync_playwright @@ -526,30 +517,18 @@ def run_bot() -> int: # noqa: C901 — orchestration, explicit branches rt["bridge"].teardown() return 3 - # Chrome env: if realtime is live on Linux, point PULSE_SOURCE at the - # virtual source so Chrome's fake mic reads the audio we generate. - chrome_env = os.environ.copy() - chrome_args = [ - "--use-fake-ui-for-media-stream", - "--disable-blink-features=AutomationControlled", - ] + chrome_args = ["--use-fake-ui-for-media-stream", "--disable-blink-features=AutomationControlled"] if not rt["enabled"]: - # v1-style fake device (silence) — we don't care about mic content - # when we're not speaking. + # Silent fake device — mic content is irrelevant when we're not speaking. chrome_args.insert(1, "--use-fake-device-for-media-stream") elif rt["bridge_info"] and rt["bridge_info"].get("platform") == "linux": - chrome_env["PULSE_SOURCE"] = rt["bridge_info"].get("device_name", "") + # Playwright's launch() takes no env: set PULSE_SOURCE on our own + # process so the child Chrome inherits the virtual source. + os.environ["PULSE_SOURCE"] = rt["bridge_info"].get("device_name", "") try: with sync_playwright() as pw: - # Playwright's launch() doesn't take env; we set PULSE_SOURCE - # via the process env before launch so the child Chrome inherits it. - for k, v in chrome_env.items(): - os.environ[k] = v - browser = pw.chromium.launch( - headless=not headed, - args=chrome_args, - ) + browser = pw.chromium.launch(headless=not cfg.headed, args=chrome_args) context_args = { "viewport": {"width": 1280, "height": 800}, "user_agent": ( @@ -558,178 +537,36 @@ def run_bot() -> int: # noqa: C901 — orchestration, explicit branches ), "permissions": ["microphone", "camera"], } - if auth_state and Path(auth_state).is_file(): - context_args["storage_state"] = auth_state + if cfg.auth_state and Path(cfg.auth_state).is_file(): + context_args["storage_state"] = cfg.auth_state context = browser.new_context(**context_args) page = context.new_page() try: - page.goto(url, wait_until="domcontentloaded", timeout=30_000) + page.goto(cfg.url, wait_until="domcontentloaded", timeout=30_000) except Exception as e: state.set(error=f"navigate failed: {e}", exited=True) return 4 - # Guest-mode: Meet shows a name field before "Ask to join". When - # we're authed, we instead see "Join now". - _try_guest_name(page, guest_name) - _click_join(page, state) - - # Install caption observer and attempt to enable captions. - try: - page.evaluate(_enable_captions_js()) + _join(page, cfg, state) + if _quiet(page.evaluate, _ENABLE_CAPTIONS_JS): state.set(captions_enabled_attempted=True) - except Exception: - pass try: page.evaluate(_CAPTION_OBSERVER_JS) except Exception as e: state.set(error=f"caption observer install failed: {e}") - # Note: in_call=False until admission is confirmed (we detect - # either the Leave button or the caption region, signalling we - # made it past the lobby). + # in_call stays False until admission is confirmed by the drain loop. state.set(captioning=True, join_attempted_at=time.time()) - - # v2 realtime: start the speaker thread reading from the - # plugin-side say queue. The thread reads JSONL lines written by - # meet_say, calls OpenAI Realtime, and streams the audio PCM to - # the virtual sink that Chrome's fake-mic is pointed at. if rt["enabled"]: - _start_realtime_speaker( - rt=rt, - out_dir=out_dir, - bridge_info=rt["bridge_info"], - api_key=realtime_api_key, - model=realtime_model, - voice=realtime_voice, - instructions=realtime_instructions, - stop_flag=stop_flag, - state=state, - ) - if rt["session"] is not None: - state.set(realtime_ready=True) + _start_realtime_speaker(rt, cfg, stop_flag, state) - # Admission + drain loop. Runs until SIGTERM, duration expiry, - # or the page detects "You were removed / you left the - # meeting". Responsible for: - # * detecting admission (Leave button visible → in_call=True) - # * timing out stuck-in-lobby (default 5 minutes) - # * draining scraped captions into the transcript - # * triggering realtime barge-in when a human speaks while - # the bot is generating audio - # * periodically flushing realtime counters into status.json - deadline = (time.time() + duration_s) if duration_s else None - lobby_deadline = time.time() + float( - os.environ.get("HERMES_MEET_LOBBY_TIMEOUT", "300") - ) - last_admission_check = 0.0 - while not stop_flag["stop"]: - now = time.time() - if deadline and now > deadline: - state.set(leave_reason="duration_expired") - break - - # Admission detection every ~3s until admitted. - if not state.in_call and (now - last_admission_check) > 3.0: - last_admission_check = now - admitted = _detect_admission(page) - if admitted: - state.set( - in_call=True, - lobby_waiting=False, - joined_at=now, - ) - elif now > lobby_deadline: - state.set( - error=( - "lobby timeout — host never admitted the bot " - f"within {int(lobby_deadline - state.join_attempted_at) if state.join_attempted_at else 0}s" - ), - leave_reason="lobby_timeout", - ) - break - elif _detect_denied(page): - state.set( - error="host denied admission", - leave_reason="denied", - ) - break - - try: - queued = page.evaluate("window.__hermesMeetDrain && window.__hermesMeetDrain()") - if isinstance(queued, list): - for entry in queued: - if not isinstance(entry, dict): - continue - speaker = str(entry.get("speaker", "")) - text = str(entry.get("text", "")) - state.record_caption(speaker=speaker, text=text) - # Barge-in: if the bot is currently generating - # audio AND a real human just spoke, cancel the - # in-flight response so we don't talk over them. - if rt["enabled"] and rt["session"] is not None: - if _looks_like_human_speaker(speaker, guest_name): - try: - cancelled = rt["session"].cancel_response() - if cancelled: - state.set(last_barge_in_at=now) - except Exception: - pass - except Exception: - # Meet reloaded or we got booted — try to detect and - # exit gracefully rather than spinning. - if page.is_closed(): - state.set(leave_reason="page_closed") - break - - # Fold the realtime session's byte/timestamp counters into - # the status file so meet_status can surface them. - if rt["session"] is not None: - state.set( - audio_bytes_out=getattr(rt["session"], "audio_bytes_out", 0), - last_audio_out_at=getattr(rt["session"], "last_audio_out_at", None), - ) - - time.sleep(1.0) - - # Try to leave cleanly — click "Leave call" button if present. - try: - page.evaluate( - "() => { const b = document.querySelector('button[aria-label*=\"eave call\"]');" - " if (b) b.click(); }" - ) - except Exception: - pass + _drain_loop(page, cfg, state, rt, stop_flag) + _quiet(page.evaluate, _LEAVE_CALL_JS) context.close() browser.close() - # v2: teardown PCM pump, speaker thread, and audio bridge. - if rt.get("pcm_pump"): - try: - rt["pcm_pump"].terminate() - rt["pcm_pump"].wait(timeout=3) - except Exception: - pass - if rt["speaker_stop"]: - try: - rt["speaker_stop"]() - except Exception: - pass - if rt["speaker_thread"] is not None: - try: - rt["speaker_thread"].join(timeout=5.0) - except Exception: - pass - if rt["session"]: - try: - rt["session"].close() - except Exception: - pass - if rt["bridge"]: - try: - rt["bridge"].teardown() - except Exception: - pass + _teardown_realtime(rt) state.set(in_call=False, captioning=False, exited=True) return 0 @@ -738,107 +575,18 @@ def run_bot() -> int: # noqa: C901 — orchestration, explicit branches return 1 -def _try_guest_name(page, guest_name: str) -> None: - """If Meet is showing a guest-name input, type *guest_name* into it.""" - try: - # Meet's guest name input has placeholder "Your name". - locator = page.locator('input[aria-label*="name" i]').first - if locator.count() and locator.is_visible(): - locator.fill(guest_name, timeout=2_000) - except Exception: - pass - - -def _detect_admission(page) -> bool: - """True if we're clearly past the lobby and in the call itself. - - Uses a JS-side probe because Meet's DOM structure varies by client - version. We check several high-signal indicators and declare admission - on the first hit: - - 1. Leave-call button is present (``aria-label`` contains "eave call"). - 2. Caption region has appeared (we installed the observer and it attached). - 3. The participant list container is visible. - - Conservative by default — returns False on any error. - """ - probe = r""" - (() => { - const leave = document.querySelector('button[aria-label*="eave call" i]'); - if (leave) return true; - if (window.__hermesMeetInstalled) { - const caps = document.querySelector( - '[role="region"][aria-label*="aption" i], ' + - 'div[jsname="YSxPC"], div[jsname="tgaKEf"]' - ); - if (caps) return true; - } - const parts = document.querySelector('[aria-label*="articipants" i]'); - if (parts) return true; - return false; - })(); - """ - try: - return bool(page.evaluate(probe)) - except Exception: - return False - - -def _detect_denied(page) -> bool: - """True when Meet is showing a 'you were denied' / 'no one admitted' page.""" - probe = r""" - (() => { - const text = document.body ? document.body.innerText || '' : ''; - // English only — matches what shows up when the host denies or - // removes a guest. - if (/You can't join this video call/i.test(text)) return true; - if (/You were removed from the meeting/i.test(text)) return true; - if (/No one responded to your request to join/i.test(text)) return true; - return false; - })(); - """ - try: - return bool(page.evaluate(probe)) - except Exception: - return False - - def _looks_like_human_speaker(speaker: str, bot_guest_name: str) -> bool: - """Whether a caption line's speaker is probably a human, not our bot echo. + """Whether a caption's speaker is probably a human rather than our own echo. - Meet attributes captions to the speaker's display name. When Chrome is - reading our fake mic, Meet still attributes captions to *our* bot name - (because the bot is the one "speaking"). We don't want those to trigger - barge-in. Anything else — real participant names — does. - - Conservative: unknown / blank speakers (common when caption scraping - falls back to raw text) do NOT trigger barge-in, because we can't tell - whether it was a human or us. + Meet attributes our fake-mic audio to the bot's own name; blank/unknown + speakers (raw-text fallback) are ambiguous, so neither triggers barge-in. """ if not speaker or not speaker.strip(): return False - spk = speaker.strip().lower() - if spk in {"unknown", "you", bot_guest_name.strip().lower()}: - return False - return True + return speaker.strip().lower() not in {"unknown", "you", bot_guest_name.strip().lower()} -def _click_join(page, state: _BotState) -> None: - """Click 'Join now' or 'Ask to join' if either button is visible. - - Flags ``lobby_waiting`` when we hit the "waiting for host to admit you" - state so the agent can surface that in status. - """ - for label in ("Join now", "Ask to join"): - try: - btn = page.get_by_role("button", name=label, exact=False).first - if btn.count() and btn.is_visible(): - btn.click(timeout=3_000) - if label == "Ask to join": - state.set(lobby_waiting=True) - break - except Exception: - continue +_DURATION_UNITS = {"h": 3600.0, "m": 60.0, "s": 1.0} def _parse_duration(raw: str) -> Optional[float]: @@ -846,14 +594,9 @@ def _parse_duration(raw: str) -> Optional[float]: if not raw: return None raw = raw.strip().lower() + mult = _DURATION_UNITS.get(raw[-1:]) try: - if raw.endswith("h"): - return float(raw[:-1]) * 3600 - if raw.endswith("m"): - return float(raw[:-1]) * 60 - if raw.endswith("s"): - return float(raw[:-1]) - return float(raw) + return float(raw[:-1]) * mult if mult else float(raw) except ValueError: return None diff --git a/plugins/google_meet/node/__init__.py b/plugins/google_meet/node/__init__.py index 338203b329..9cbccaf438 100644 --- a/plugins/google_meet/node/__init__.py +++ b/plugins/google_meet/node/__init__.py @@ -1,54 +1,25 @@ -"""Remote 'node host' primitive for the google_meet plugin. +"""Remote 'node host' primitive: run the Meet bot on a different machine than the gateway. -Lets the Meet bot (Playwright + Chrome) run on a different machine than -the hermes-agent gateway. The gateway speaks a small JSON-over-WebSocket -RPC protocol to the remote node; the node wraps the existing -``plugins.google_meet.process_manager`` API. + gateway (Linux) ── ws://mac.local:18789 ──▶ node server (Mac) → process_manager → meet_bot -Topology --------- - gateway (Linux) ── ws://mac.local:18789 ──▶ node server (Mac) - └─ process_manager - └─ meet_bot (Playwright) +Why: Google sign-in + Chrome profile live on the user's laptop; running the +bot there reuses that profile without shipping credentials to the server. -Why: Google sign-in + Chrome profile live on the user's laptop. Running -the bot there reuses that profile without shipping credentials to the -server. - -Public surface --------------- - NodeClient — gateway-side RPC client (short-lived sync WS per call) - NodeServer — long-running server that hosts the bot - NodeRegistry — local JSON registry of approved nodes (name → url+token) - protocol — message envelope helpers (make_request, encode, decode, ...) +NodeClient (gateway-side RPC), NodeServer (hosts the bot), NodeRegistry +(approved nodes: name → url+token), protocol (envelope helpers). """ from __future__ import annotations from plugins.google_meet.node import protocol from plugins.google_meet.node.client import NodeClient -from plugins.google_meet.node.protocol import ( - VALID_REQUEST_TYPES, - decode, - encode, - make_error, - make_request, - make_response, - validate_request, +from plugins.google_meet.node.protocol import ( # noqa: F401 + VALID_REQUEST_TYPES, decode, encode, make_error, make_request, make_response, validate_request, ) from plugins.google_meet.node.registry import NodeRegistry from plugins.google_meet.node.server import NodeServer __all__ = [ - "NodeClient", - "NodeServer", - "NodeRegistry", - "protocol", - "make_request", - "make_response", - "make_error", - "encode", - "decode", - "validate_request", - "VALID_REQUEST_TYPES", + "NodeClient", "NodeServer", "NodeRegistry", "protocol", "make_request", "make_response", + "make_error", "encode", "decode", "validate_request", "VALID_REQUEST_TYPES", ] diff --git a/plugins/google_meet/node/cli.py b/plugins/google_meet/node/cli.py index ac2b08ac58..6553797ba4 100644 --- a/plugins/google_meet/node/cli.py +++ b/plugins/google_meet/node/cli.py @@ -1,9 +1,4 @@ -"""`hermes meet node ...` subcommand tree. - -Wired into the existing ``hermes meet`` parser by the plugin's top-level -CLI. This module only defines the subparsers and their dispatch — it -does not mutate the existing cli.py. -""" +"""`hermes meet node ...` subcommand tree, wired under the ``hermes meet`` parser.""" from __future__ import annotations @@ -11,7 +6,6 @@ import argparse import asyncio import json import sys -from typing import Any from plugins.google_meet.node.client import NodeClient from plugins.google_meet.node.registry import NodeRegistry @@ -19,107 +13,100 @@ from plugins.google_meet.node.server import NodeServer def register_cli(subparser: argparse.ArgumentParser) -> None: - """Add ``run / list / approve / remove / status / ping`` subparsers. - - *subparser* is the ``hermes meet node`` argparse object — typically - the result of ``meet_parser.add_parser('node', ...)``. - """ + """Add ``run / list / approve / remove / status / ping`` subparsers to the ``node`` parser.""" sp = subparser.add_subparsers(dest="node_cmd", required=True) run = sp.add_parser("run", help="Start a node server on this machine.") run.add_argument("--host", default="0.0.0.0") run.add_argument("--port", type=int, default=18789) run.add_argument("--display-name", default="hermes-meet-node") - run.set_defaults(func=node_command) - lst = sp.add_parser("list", help="List approved remote nodes.") - lst.set_defaults(func=node_command) + sp.add_parser("list", help="List approved remote nodes.") app = sp.add_parser("approve", help="Register a remote node on the gateway.") - app.add_argument("name") - app.add_argument("url") - app.add_argument("token") - app.set_defaults(func=node_command) + for arg in ("name", "url", "token"): + app.add_argument(arg) - rm = sp.add_parser("remove", help="Forget a registered node.") - rm.add_argument("name") - rm.set_defaults(func=node_command) + for name, help_ in (("remove", "Forget a registered node."), + ("status", "Ping a registered node."), + ("ping", "Alias for status.")): + sp.add_parser(name, help=help_).add_argument("name") - st = sp.add_parser("status", help="Ping a registered node.") - st.add_argument("name") - st.set_defaults(func=node_command) + for p in sp.choices.values(): + p.set_defaults(func=node_command) - pg = sp.add_parser("ping", help="Alias for status.") - pg.add_argument("name") - pg.set_defaults(func=node_command) + +def _cmd_run(args: argparse.Namespace, reg: NodeRegistry) -> int: + server = NodeServer(host=args.host, port=args.port, display_name=args.display_name) + token = server.ensure_token() + print(f"[meet-node] display_name={server.display_name}\n" + f"[meet-node] listening on ws://{args.host}:{args.port}\n" + f"[meet-node] token (copy to gateway): {token}\n" + "[meet-node] approve with:\n" + f" hermes meet node approve ws://:{args.port} {token}") + try: + asyncio.run(server.serve()) + except KeyboardInterrupt: + pass + except RuntimeError as exc: + print(f"[meet-node] error: {exc}", file=sys.stderr) + return 2 + return 0 + + +def _cmd_list(args: argparse.Namespace, reg: NodeRegistry) -> int: + nodes = reg.list_all() + if not nodes: + print("no nodes registered") + for n in nodes: + print(f"{n['name']}\t{n['url']}\ttoken={n['token'][:6]}…") + return 0 + + +def _cmd_approve(args: argparse.Namespace, reg: NodeRegistry) -> int: + reg.add(args.name, args.url, args.token) + print(f"approved node {args.name!r} at {args.url}") + return 0 + + +def _cmd_remove(args: argparse.Namespace, reg: NodeRegistry) -> int: + ok = reg.remove(args.name) + print(f"removed {args.name!r}" if ok else f"no such node: {args.name!r}") + return 0 if ok else 1 + + +def _cmd_ping(args: argparse.Namespace, reg: NodeRegistry) -> int: + entry = reg.get(args.name) + if entry is None: + print(f"no such node: {args.name!r}", file=sys.stderr) + return 1 + try: + result = NodeClient(entry["url"], entry["token"]).ping() + except Exception as exc: # noqa: BLE001 — surface any connection error + print(json.dumps({"ok": False, "error": str(exc)})) + return 1 + if not isinstance(result, dict): + result = {"result": result} + print(json.dumps({"ok": True, "node": args.name, **result})) + return 0 + + +_COMMANDS = { + "run": _cmd_run, + "list": _cmd_list, + "approve": _cmd_approve, + "remove": _cmd_remove, + "status": _cmd_ping, + "ping": _cmd_ping, +} def node_command(args: argparse.Namespace) -> int: - """Dispatch for ``hermes meet node ...``. - - Returns a process exit code. Side-effects print to stdout/stderr. - """ + """Dispatch for ``hermes meet node ...``; returns a process exit code.""" cmd = getattr(args, "node_cmd", None) - - if cmd == "run": - server = NodeServer( - host=args.host, - port=args.port, - display_name=args.display_name, - ) - token = server.ensure_token() - print(f"[meet-node] display_name={server.display_name}") - print(f"[meet-node] listening on ws://{args.host}:{args.port}") - print(f"[meet-node] token (copy to gateway): {token}") - print("[meet-node] approve with:") - print(f" hermes meet node approve ws://:{args.port} {token}") - try: - asyncio.run(server.serve()) - except KeyboardInterrupt: - return 0 - except RuntimeError as exc: - print(f"[meet-node] error: {exc}", file=sys.stderr) - return 2 - return 0 - - reg = NodeRegistry() - - if cmd == "list": - nodes = reg.list_all() - if not nodes: - print("no nodes registered") - return 0 - for n in nodes: - print(f"{n['name']}\t{n['url']}\ttoken={n['token'][:6]}…") - return 0 - - if cmd == "approve": - reg.add(args.name, args.url, args.token) - print(f"approved node {args.name!r} at {args.url}") - return 0 - - if cmd == "remove": - ok = reg.remove(args.name) - print(f"removed {args.name!r}" if ok else f"no such node: {args.name!r}") - return 0 if ok else 1 - - if cmd in {"status", "ping"}: - entry = reg.get(args.name) - if entry is None: - print(f"no such node: {args.name!r}", file=sys.stderr) - return 1 - client = NodeClient(entry["url"], entry["token"]) - try: - result = client.ping() - except Exception as exc: # noqa: BLE001 — surface any connection error - print(json.dumps({"ok": False, "error": str(exc)})) - return 1 - print(json.dumps({"ok": True, "node": args.name, **_coerce_dict(result)})) - return 0 - - print(f"unknown node command: {cmd!r}", file=sys.stderr) - return 2 - - -def _coerce_dict(value: Any) -> dict: - return value if isinstance(value, dict) else {"result": value} + handler = _COMMANDS.get(cmd or "") + if handler is None: + print(f"unknown node command: {cmd!r}", file=sys.stderr) + return 2 + # ``run`` never touches the registry; constructing it is side-effect free. + return handler(args, NodeRegistry()) diff --git a/plugins/google_meet/node/client.py b/plugins/google_meet/node/client.py index 1965333c0b..840f7253ad 100644 --- a/plugins/google_meet/node/client.py +++ b/plugins/google_meet/node/client.py @@ -1,12 +1,8 @@ """Gateway-side RPC client for a remote meet node. -Each call opens a short-lived synchronous WebSocket to the node, sends -exactly one request, reads exactly one response, and closes. This keeps -the client trivial to use from non-async tool handlers and avoids -maintaining persistent connection state across agent turns. - -The ``websockets`` package is an optional dep — we import it lazily so -plugin load doesn't require it. +One short-lived sync WebSocket per call (send one request, read one +response, close) so non-async tool handlers need no persistent connection +state across agent turns. ``websockets`` is optional and imported lazily. """ from __future__ import annotations @@ -28,14 +24,8 @@ class NodeClient: self.token = token self.timeout = float(timeout) - # ----- core RPC ----------------------------------------------------- - def _rpc(self, type: str, payload: Dict[str, Any]) -> Dict[str, Any]: - """Send one request, return the response payload dict. - - Raises RuntimeError when the server sends an ``error`` envelope - or the response id doesn't match. - """ + """Send one request, return its payload dict; RuntimeError on error envelope / id mismatch.""" try: from websockets.sync.client import connect # type: ignore except ImportError as exc: @@ -45,31 +35,19 @@ class NodeClient: ) from exc req = _proto.make_request(type, self.token, payload) - raw_out = _proto.encode(req) - - with connect(self.url, open_timeout=self.timeout, - close_timeout=self.timeout) as ws: - ws.send(raw_out) - raw_in = ws.recv(timeout=self.timeout) - - if isinstance(raw_in, (bytes, bytearray)): - raw_in = raw_in.decode("utf-8") - resp = _proto.decode(raw_in) + with connect(self.url, open_timeout=self.timeout, close_timeout=self.timeout) as ws: + ws.send(_proto.encode(req)) + resp = _proto.decode(ws.recv(timeout=self.timeout)) if resp.get("type") == "error": raise RuntimeError(f"node error: {resp.get('error', '')}") if resp.get("id") != req["id"]: - raise RuntimeError( - f"response id mismatch: sent {req['id']}, got {resp.get('id')!r}" - ) + raise RuntimeError(f"response id mismatch: sent {req['id']}, got {resp.get('id')!r}") payload_out = resp.get("payload") - if not isinstance(payload_out, dict): - # Ping returns {"type": "pong", "payload": {...}} — still a dict. + if not isinstance(payload_out, dict): # pong envelopes also carry a dict payload raise RuntimeError("response missing payload dict") return payload_out - # ----- convenience methods ----------------------------------------- - def start_bot( self, url: str, @@ -78,12 +56,7 @@ class NodeClient: headed: bool = False, mode: str = "transcribe", ) -> Dict[str, Any]: - payload: Dict[str, Any] = { - "url": url, - "guest_name": guest_name, - "headed": bool(headed), - "mode": mode, - } + payload: Dict[str, Any] = {"url": url, "guest_name": guest_name, "headed": bool(headed), "mode": mode} if duration is not None: payload["duration"] = duration return self._rpc("start_bot", payload) @@ -95,10 +68,7 @@ class NodeClient: return self._rpc("status", {}) def transcript(self, last: Optional[int] = None) -> Dict[str, Any]: - payload: Dict[str, Any] = {} - if last is not None: - payload["last"] = int(last) - return self._rpc("transcript", payload) + return self._rpc("transcript", {} if last is None else {"last": int(last)}) def say(self, text: str) -> Dict[str, Any]: return self._rpc("say", {"text": str(text)}) diff --git a/plugins/google_meet/node/protocol.py b/plugins/google_meet/node/protocol.py index 8794d8a533..5e32643739 100644 --- a/plugins/google_meet/node/protocol.py +++ b/plugins/google_meet/node/protocol.py @@ -3,12 +3,12 @@ Everything is a JSON object with the same envelope shape: Request: {"type": , "id": , "token": , "payload": } - Response: {"type": "_res", "id": , "payload": } + Response: {"type": "response", "id": , "payload": } Error: {"type": "error", "id": , "error": } -Requests must carry the shared bearer token (set up via -``hermes meet node approve`` on the gateway and read off disk on the -server). Mismatched tokens are rejected before dispatch. +Requests must carry the shared bearer token (set up via ``hermes meet node +approve`` on the gateway and read off disk on the server). Mismatched tokens +are rejected before dispatch. """ from __future__ import annotations @@ -18,28 +18,16 @@ import uuid from typing import Any, Dict, Tuple -VALID_REQUEST_TYPES = frozenset({ - "start_bot", - "stop", - "status", - "transcript", - "say", - "ping", -}) +VALID_REQUEST_TYPES = frozenset({"start_bot", "stop", "status", "transcript", "say", "ping"}) -def make_request( - type: str, - token: str, - payload: Dict[str, Any], - req_id: str | None = None, -) -> Dict[str, Any]: - """Construct a request envelope. +def _nonempty_str(value: Any) -> bool: + return isinstance(value, str) and bool(value) - ``req_id`` is auto-generated (uuid4 hex) when not supplied so callers - can correlate async responses. - """ - if not isinstance(type, str) or not type: + +def make_request(type: str, token: str, payload: Dict[str, Any], req_id: str | None = None) -> Dict[str, Any]: + """Construct a request envelope; ``req_id`` defaults to a uuid4 hex.""" + if not _nonempty_str(type): raise ValueError("type must be a non-empty string") if type not in VALID_REQUEST_TYPES: raise ValueError(f"unknown request type: {type!r}") @@ -47,22 +35,11 @@ def make_request( raise ValueError("token must be a string") if not isinstance(payload, dict): raise ValueError("payload must be a dict") - return { - "type": type, - "id": req_id or uuid.uuid4().hex, - "token": token, - "payload": payload, - } + return {"type": type, "id": req_id or uuid.uuid4().hex, "token": token, "payload": payload} def make_response(req_id: str, payload: Dict[str, Any]) -> Dict[str, Any]: - """Build a success response. The caller supplies the *request* type; - we suffix it with ``_res`` so clients can assert they got the right - reply. - - For simplicity we don't require the type here — clients usually just - key off ``id``. But we still emit a generic ``*_res`` envelope. - """ + """Build a success envelope; clients correlate replies by ``id``, not type.""" if not isinstance(payload, dict): raise ValueError("payload must be a dict") return {"type": "response", "id": req_id, "payload": payload} @@ -77,48 +54,42 @@ def encode(msg: Dict[str, Any]) -> str: return json.dumps(msg, separators=(",", ":"), ensure_ascii=False) -def decode(raw: str) -> Dict[str, Any]: - """Parse a JSON envelope, raising ValueError on anything malformed. +def decode(raw) -> Dict[str, Any]: + """Parse a JSON envelope (object with string ``type`` + ``id``), raising ValueError otherwise. - Minimal type validation: must be an object, must contain ``type`` and - ``id``. Heavier validation (token match, payload shape) happens in - :func:`validate_request` on the server side. + Accepts ``str`` or UTF-8 ``bytes``. Token match and payload shape are + checked server-side in :func:`validate_request`. """ + if isinstance(raw, (bytes, bytearray)): + raw = raw.decode("utf-8") try: obj = json.loads(raw) except (TypeError, json.JSONDecodeError) as exc: raise ValueError(f"malformed JSON: {exc}") from exc if not isinstance(obj, dict): raise ValueError("envelope must be a JSON object") - if "type" not in obj or not isinstance(obj["type"], str): - raise ValueError("envelope missing string 'type'") - if "id" not in obj or not isinstance(obj["id"], str): - raise ValueError("envelope missing string 'id'") + for key in ("type", "id"): + if not isinstance(obj.get(key), str): + raise ValueError(f"envelope missing string '{key}'") return obj def validate_request(msg: Dict[str, Any], expected_token: str) -> Tuple[bool, str]: - """Check a decoded request against the server's shared token. - - Returns ``(True, "")`` when the envelope is acceptable or - ``(False, )`` otherwise. Reason strings are safe to surface - back to the client in an error envelope. - """ + """Return ``(True, "")`` or ``(False, )``; reasons are safe to send back to the client.""" if not isinstance(msg, dict): return False, "envelope must be a dict" t = msg.get("type") - if not isinstance(t, str) or not t: + if not _nonempty_str(t): return False, "missing or non-string 'type'" if t not in VALID_REQUEST_TYPES: return False, f"unknown request type: {t!r}" - if not isinstance(msg.get("id"), str) or not msg.get("id"): + if not _nonempty_str(msg.get("id")): return False, "missing or non-string 'id'" token = msg.get("token") - if not isinstance(token, str) or not token: + if not _nonempty_str(token): return False, "missing token" if token != expected_token: return False, "token mismatch" - payload = msg.get("payload") - if not isinstance(payload, dict): + if not isinstance(msg.get("payload"), dict): return False, "payload must be a dict" return True, "" diff --git a/plugins/google_meet/node/registry.py b/plugins/google_meet/node/registry.py index 9be8575562..4755f627e4 100644 --- a/plugins/google_meet/node/registry.py +++ b/plugins/google_meet/node/registry.py @@ -1,112 +1,73 @@ """Local JSON registry of approved remote meet nodes. -Lives at ``$HERMES_HOME/workspace/meetings/nodes.json``. The gateway -consults it to resolve a ``chrome_node`` name to a ``(url, token)`` pair -before opening a WebSocket to the remote bot host. +``$HERMES_HOME/workspace/meetings/nodes.json``:: -Schema ------- - { - "nodes": { - "": { - "url": "ws://host:port", - "token": "...", - "added_at": - } - } - } + {"nodes": {"": {"url": "ws://host:port", "token": "...", "added_at": }}} """ from __future__ import annotations -import json import time from pathlib import Path from typing import Any, Dict, List, Optional from hermes_constants import get_hermes_home +from plugins.google_meet._jsonfile import read_json, write_json_atomic + def _default_path() -> Path: return Path(get_hermes_home()) / "workspace" / "meetings" / "nodes.json" class NodeRegistry: - """Simple file-backed registry. Not concurrent-safe across processes - — single writer assumed (the gateway CLI).""" + """File-backed registry; single writer assumed (the gateway CLI).""" def __init__(self, path: Optional[Path] = None) -> None: self.path = Path(path) if path is not None else _default_path() # ----- storage ------------------------------------------------------ - def _load(self) -> Dict[str, Any]: - if not self.path.is_file(): - return {"nodes": {}} - try: - data = json.loads(self.path.read_text(encoding="utf-8")) - except (OSError, json.JSONDecodeError): - return {"nodes": {}} - if not isinstance(data, dict) or not isinstance(data.get("nodes"), dict): - return {"nodes": {}} - return data + def _load(self) -> Dict[str, Dict[str, Any]]: + """The ``nodes`` map (name → entry); empty when the file is missing or malformed.""" + data = read_json(self.path) + nodes = data.get("nodes") if isinstance(data, dict) else None + return nodes if isinstance(nodes, dict) else {} - def _save(self, data: Dict[str, Any]) -> None: - self.path.parent.mkdir(parents=True, exist_ok=True) - tmp = self.path.with_suffix(".json.tmp") - tmp.write_text(json.dumps(data, indent=2), encoding="utf-8") - tmp.replace(self.path) + def _save(self, nodes: Dict[str, Dict[str, Any]]) -> None: + write_json_atomic(self.path, {"nodes": nodes}) # ----- public API --------------------------------------------------- def get(self, name: str) -> Optional[Dict[str, Any]]: - data = self._load() - entry = data["nodes"].get(name) - if entry is None: - return None - return {"name": name, **entry} + entry = self._load().get(name) + return None if entry is None else {"name": name, **entry} def add(self, name: str, url: str, token: str) -> None: - if not isinstance(name, str) or not name: - raise ValueError("node name must be a non-empty string") - if not isinstance(url, str) or not url: - raise ValueError("url must be a non-empty string") - if not isinstance(token, str) or not token: - raise ValueError("token must be a non-empty string") - data = self._load() - data["nodes"][name] = { - "url": url, - "token": token, - "added_at": time.time(), - } - self._save(data) + for label, value in (("node name", name), ("url", url), ("token", token)): + if not isinstance(value, str) or not value: + raise ValueError(f"{label} must be a non-empty string") + nodes = self._load() + nodes[name] = {"url": url, "token": token, "added_at": time.time()} + self._save(nodes) def remove(self, name: str) -> bool: - data = self._load() - if name in data["nodes"]: - del data["nodes"][name] - self._save(data) - return True - return False + nodes = self._load() + if name not in nodes: + return False + del nodes[name] + self._save(nodes) + return True def list_all(self) -> List[Dict[str, Any]]: - data = self._load() - out: List[Dict[str, Any]] = [] - for name, entry in sorted(data["nodes"].items()): - out.append({"name": name, **entry}) - return out + return [{"name": name, **entry} for name, entry in sorted(self._load().items())] def resolve(self, chrome_node: Optional[str]) -> Optional[Dict[str, Any]]: - """Resolve a node name to its entry. + """Named node's entry, or — when ``chrome_node`` is falsy — the sole registered node. - If ``chrome_node`` is provided, return that named node (or None). - If ``chrome_node`` is None, return the sole registered node when - exactly one is registered; otherwise return None (ambiguous or - empty). + None when the name is unknown or when zero / several nodes are registered (ambiguous). """ if chrome_node: return self.get(chrome_node) nodes = self.list_all() - if len(nodes) == 1: - return nodes[0] - return None + return nodes[0] if len(nodes) == 1 else None diff --git a/plugins/google_meet/node/server.py b/plugins/google_meet/node/server.py index cff01d265f..632e2ba03e 100644 --- a/plugins/google_meet/node/server.py +++ b/plugins/google_meet/node/server.py @@ -1,29 +1,20 @@ -"""Remote node server. +"""Remote node server — hosts the Meet bot on another machine (e.g. the user's Mac). -Runs on the machine that will host the Meet bot (typically the user's -Mac laptop with a signed-in Chrome). Exposes a WebSocket endpoint that -accepts signed RPC requests and dispatches them to the existing -``plugins.google_meet.process_manager`` module. +Exposes a WebSocket endpoint that accepts token-signed RPC requests and +dispatches them to ``plugins.google_meet.process_manager``. Launched by +``hermes meet node run``. -Launched by ``hermes meet node run``. +Token: 32 hex chars minted on first boot and persisted at +``$HERMES_HOME/workspace/meetings/node_token.json`` so previously-approved +gateways survive restarts. The operator copies it to the gateway via +``hermes meet node approve ``. -Token handling --------------- -On first boot we mint 32 hex chars of entropy and persist them at -``$HERMES_HOME/workspace/meetings/node_token.json``. Subsequent boots -reuse the same token so previously-approved gateways don't need to be -re-paired. The operator copies this token out-of-band to the gateway -via ``hermes meet node approve ``. - -Dependencies ------------- -``websockets`` is an optional dep. We import it lazily inside -:meth:`serve` so installing the plugin doesn't require it unless you -actually host a node. +``websockets`` is optional and imported lazily inside :meth:`serve`. """ from __future__ import annotations +import asyncio import json import secrets import time @@ -31,11 +22,51 @@ from pathlib import Path from typing import Any, Dict, Optional from hermes_constants import get_hermes_home +from plugins.google_meet._jsonfile import read_json, write_json_atomic from plugins.google_meet.node import protocol as _proto +_START_BOT_KEYS = ("url", "guest_name", "duration", "headed", "auth_state", "session_id", "out_dir") -def _default_token_path() -> Path: - return Path(get_hermes_home()) / "workspace" / "meetings" / "node_token.json" + +class _RpcError(Exception): + """Handler-level protocol error; sent verbatim as an error envelope.""" + + +def _rpc_start_bot(payload: Dict[str, Any], pm) -> Dict[str, Any]: + # Whitelist kwargs we pass through to pm.start. + kwargs = {k: payload[k] for k in _START_BOT_KEYS if k in payload} + if "url" not in kwargs: + raise _RpcError("missing 'url' in payload") + return pm.start(**kwargs) + + +def _rpc_say(payload: Dict[str, Any], pm) -> Dict[str, Any]: + # Appends to say_queue.jsonl inside the active meeting's out_dir; the + # bot-side consumer only exists in realtime mode, so ok=True here means + # "enqueued", not "spoken". + text = payload.get("text", "") + active = pm._read_active() + enqueued = False + if active and active.get("out_dir"): + queue = Path(active["out_dir"]) / "say_queue.jsonl" + try: + queue.parent.mkdir(parents=True, exist_ok=True) + with queue.open("a", encoding="utf-8") as fh: + fh.write(json.dumps({"text": text, "ts": time.time()}) + "\n") + enqueued = True + except OSError: + pass + return {"ok": True, "enqueued": enqueued, "text": text} + + +# request type → fn(payload, pm) returning the response payload. +_RPC = { + "start_bot": _rpc_start_bot, + "stop": lambda p, pm: pm.stop(reason=p.get("reason", "requested")), + "status": lambda p, pm: pm.status(), + "transcript": lambda p, pm: pm.transcript(last=p.get("last")), + "say": _rpc_say, +} class NodeServer: @@ -51,129 +82,54 @@ class NodeServer: self.host = host self.port = port self.display_name = display_name - self.token_path = Path(token_path) if token_path is not None else _default_token_path() + self.token_path = Path(token_path) if token_path is not None else ( + Path(get_hermes_home()) / "workspace" / "meetings" / "node_token.json" + ) self._token: Optional[str] = None - # ----- token management -------------------------------------------- - def ensure_token(self) -> str: """Return the persisted shared secret, generating one on first use.""" if self._token: return self._token - if self.token_path.is_file(): - try: - data = json.loads(self.token_path.read_text(encoding="utf-8")) - tok = data.get("token") - if isinstance(tok, str) and tok: - self._token = tok - return tok - except (OSError, json.JSONDecodeError): - pass - tok = secrets.token_hex(16) # 32 hex chars - self.token_path.parent.mkdir(parents=True, exist_ok=True) - tmp = self.token_path.with_suffix(".json.tmp") - tmp.write_text( - json.dumps({"token": tok, "generated_at": time.time()}, indent=2), - encoding="utf-8", - ) - # Restrict to owner-read-write only — the token grants full RPC - # access to the meet bot (start, transcribe, speak in meetings). - try: - tmp.chmod(0o600) - except (OSError, NotImplementedError): - # Best-effort on non-POSIX filesystems; mode is set on POSIX. - pass - tmp.replace(self.token_path) + data = read_json(self.token_path) + tok = data.get("token") if isinstance(data, dict) else None + if not (isinstance(tok, str) and tok): + tok = secrets.token_hex(16) # 32 hex chars + # Owner-only: the token grants full RPC access to the meet bot. + write_json_atomic(self.token_path, {"token": tok, "generated_at": time.time()}, mode=0o600) self._token = tok return tok - def get_token(self) -> str: - """Alias for :meth:`ensure_token`; does not mutate on subsequent calls.""" - return self.ensure_token() - - # ----- dispatch ----------------------------------------------------- - async def _handle_request(self, msg: Dict[str, Any]) -> Dict[str, Any]: - """Validate + dispatch a single decoded request envelope. + """Validate + dispatch one decoded request; always returns an envelope, never raises. - Always returns a response envelope (success or error); never - raises. Errors from inside the process_manager are wrapped into - the response payload's ``ok``/``error`` keys (which pm already - does) rather than being re-encoded as error envelopes — the - envelope-level error channel is reserved for auth / protocol - failures. + The envelope-level ``error`` channel is reserved for auth/protocol + failures and pm crashes; pm's own ``ok``/``error`` results travel + inside a normal response payload. """ - expected = self.ensure_token() - ok, reason = _proto.validate_request(msg, expected) + ok, reason = _proto.validate_request(msg, self.ensure_token()) if not ok: return _proto.make_error(str(msg.get("id") or ""), reason) - req_id = msg["id"] - t = msg["type"] - payload = msg["payload"] + req_id, t = msg["id"], msg["type"] + if t == "ping": + return {"type": "pong", "id": req_id, "payload": {"display_name": self.display_name, "ts": time.time()}} + handler = _RPC.get(t) + if handler is None: + return _proto.make_error(req_id, f"unhandled type: {t!r}") # Import lazily so test mocks can monkeypatch freely. from plugins.google_meet import process_manager as pm try: - if t == "ping": - return {"type": "pong", "id": req_id, - "payload": {"display_name": self.display_name, - "ts": time.time()}} - if t == "start_bot": - # Whitelist kwargs we pass through to pm.start. - kwargs = { - k: payload[k] - for k in ("url", "guest_name", "duration", "headed", - "auth_state", "session_id", "out_dir") - if k in payload - } - if "url" not in kwargs: - return _proto.make_error(req_id, "missing 'url' in payload") - result = pm.start(**kwargs) - return _proto.make_response(req_id, result) - if t == "stop": - reason_arg = payload.get("reason", "requested") - result = pm.stop(reason=reason_arg) - return _proto.make_response(req_id, result) - if t == "status": - return _proto.make_response(req_id, pm.status()) - if t == "transcript": - last = payload.get("last") - result = pm.transcript(last=last) - return _proto.make_response(req_id, result) - if t == "say": - # v2 wiring: enqueue into say_queue.jsonl inside the - # active meeting's out_dir when present. The bot-side - # consumer is v3+ (for v1 this is a stub returning ok). - text = payload.get("text", "") - active = pm._read_active() # type: ignore[attr-defined] - enqueued = False - if active and active.get("out_dir"): - queue = Path(active["out_dir"]) / "say_queue.jsonl" - try: - queue.parent.mkdir(parents=True, exist_ok=True) - with queue.open("a", encoding="utf-8") as fh: - fh.write(json.dumps({"text": text, "ts": time.time()}) + "\n") - enqueued = True - except OSError: - enqueued = False - return _proto.make_response( - req_id, - {"ok": True, "enqueued": enqueued, "text": text}, - ) + return _proto.make_response(req_id, handler(msg["payload"], pm)) + except _RpcError as exc: + return _proto.make_error(req_id, str(exc)) except Exception as exc: # noqa: BLE001 — surface any pm crash to client return _proto.make_error(req_id, f"{type(exc).__name__}: {exc}") - return _proto.make_error(req_id, f"unhandled type: {t!r}") - - # ----- server loop -------------------------------------------------- - async def serve(self) -> None: - """Run the WebSocket server until cancelled. - - Blocks forever. Callers typically wrap this in ``asyncio.run``. - """ + """Run the WebSocket server until cancelled (wrap in ``asyncio.run``).""" try: import websockets # type: ignore except ImportError as exc: @@ -187,14 +143,11 @@ class NodeServer: async def _handler(ws): async for raw in ws: try: - msg = _proto.decode(raw if isinstance(raw, str) else raw.decode("utf-8")) + msg = _proto.decode(raw) except ValueError as exc: await ws.send(_proto.encode(_proto.make_error("", f"decode: {exc}"))) continue - reply = await self._handle_request(msg) - await ws.send(_proto.encode(reply)) + await ws.send(_proto.encode(await self._handle_request(msg))) async with websockets.serve(_handler, self.host, self.port): - # Run until cancelled. - import asyncio - await asyncio.Future() + await asyncio.Future() # run until cancelled diff --git a/plugins/google_meet/process_manager.py b/plugins/google_meet/process_manager.py index 4b2576efb4..b28a80b2f1 100644 --- a/plugins/google_meet/process_manager.py +++ b/plugins/google_meet/process_manager.py @@ -1,12 +1,17 @@ """Subprocess lifecycle manager for the google_meet bot. -Single active meeting at a time. Stores the running pid + out_dir in a -session-scoped state file under ``$HERMES_HOME/workspace/meetings/.active.json`` -so tool calls across turns can find the bot, and ``on_session_end`` can clean -it up. +Single active meeting at a time, recorded in +``$HERMES_HOME/workspace/meetings/.active.json`` so tool calls across turns +(and ``on_session_end``) can find the bot. The bot is a detached subprocess: +we hold no fds on it and communicate via files only, so the agent loop can't +block on it. -The bot runs as a detached subprocess — we don't hold file descriptors open, -so the parent agent loop can't block on it. We communicate via files only. +Layout under ``workspace/meetings/``:: + + .active.json {"pid", "meeting_id", "out_dir", "url", "started_at", + "session_id", "log_path", "mode"} + /status.json live bot state (written by the bot) + /transcript.txt scraped captions """ from __future__ import annotations @@ -22,64 +27,42 @@ from typing import Any, Dict, Optional from hermes_constants import get_hermes_home -# File + directory layout (under $HERMES_HOME): -# -# workspace/meetings/ -# .active.json # pointer to current session's bot -# / -# status.json # live bot state (written by bot each tick) -# transcript.txt # scraped captions -# -# .active.json holds: -# {"pid": 12345, "meeting_id": "abc-defg-hij", "out_dir": "...", -# "url": "https://meet.google.com/...", "started_at": 1714159200.0, -# "session_id": "optional"} +from plugins.google_meet._jsonfile import read_json, write_json_atomic def _root() -> Path: return Path(get_hermes_home()) / "workspace" / "meetings" -def _active_file() -> Path: - return _root() / ".active.json" - - def _read_active() -> Optional[Dict[str, Any]]: - p = _active_file() - if not p.is_file(): - return None - try: - return json.loads(p.read_text(encoding="utf-8")) - except Exception: - return None + return read_json(_root() / ".active.json") def _write_active(data: Dict[str, Any]) -> None: - p = _active_file() - p.parent.mkdir(parents=True, exist_ok=True) - tmp = p.with_suffix(".json.tmp") - tmp.write_text(json.dumps(data, indent=2), encoding="utf-8") - tmp.replace(p) - - -def _clear_active() -> None: - try: - _active_file().unlink() - except FileNotFoundError: - pass + write_json_atomic(_root() / ".active.json", data) def _pid_alive(pid: int) -> bool: - # ``os.kill(pid, 0)`` is NOT a no-op on Windows (bpo-14484) — it - # routes through GenerateConsoleCtrlEvent and can kill the target. - # Use the cross-platform existence check. + # Not ``os.kill(pid, 0)``: on Windows that routes through + # GenerateConsoleCtrlEvent and can kill the target (bpo-14484). from gateway.status import _pid_exists - return _pid_exists(pid) + return bool(pid) and _pid_exists(pid) -# --------------------------------------------------------------------------- -# Public API — used by tool handlers + CLI -# --------------------------------------------------------------------------- +def _active_pid() -> int: + active = _read_active() + return int(active.get("pid", 0)) if active else 0 + + +def _kill(pid: int, sig) -> None: + try: + os.kill(pid, sig) + except ProcessLookupError: + pass + + +_NO_ACTIVE = {"ok": False, "reason": "no active meeting"} + def start( url: str, @@ -96,108 +79,64 @@ def start( realtime_instructions: Optional[str] = None, realtime_api_key: Optional[str] = None, ) -> Dict[str, Any]: - """Spawn the meet_bot subprocess for *url*. - - If a bot is already running for this hermes install, leave it first — - we enforce single-active-meeting semantics. - - Returns a dict summarizing the started bot. - """ + """Spawn the meet_bot subprocess for *url*, stopping any running bot first + (single-active-meeting semantics). Returns a dict summarizing the bot.""" from plugins.google_meet.meet_bot import _is_safe_meet_url, _meeting_id_from_url if not _is_safe_meet_url(url): - return { - "ok": False, - "error": ( - "refusing: only https://meet.google.com/ URLs are allowed. " - "got: " + repr(url) - ), - } + return {"ok": False, "error": "refusing: only https://meet.google.com/ URLs are allowed. got: " + repr(url)} - existing = _read_active() - if existing and _pid_alive(int(existing.get("pid", 0))): + if _pid_alive(_active_pid()): stop(reason="replaced by new meet_join") meeting_id = _meeting_id_from_url(url) out = out_dir or (_root() / meeting_id) out.mkdir(parents=True, exist_ok=True) - # Wipe any stale transcript/status files from a previous run of this - # meeting id so polling isn't confused. + # Wipe stale files from a previous run of this meeting id so polling isn't confused. for name in ("transcript.txt", "status.json"): - f = out / name - if f.exists(): - try: - f.unlink() - except OSError: - pass - - env = os.environ.copy() - env["HERMES_MEET_URL"] = url - env["HERMES_MEET_OUT_DIR"] = str(out) - env["HERMES_MEET_GUEST_NAME"] = guest_name - if headed: - env["HERMES_MEET_HEADED"] = "1" - if auth_state: - env["HERMES_MEET_AUTH_STATE"] = auth_state - if duration: - env["HERMES_MEET_DURATION"] = duration - # v2: realtime mode + passthroughs. The bot defaults to transcribe - # mode if HERMES_MEET_MODE isn't set, matching v1 behavior. - if mode: - env["HERMES_MEET_MODE"] = mode - if realtime_model: - env["HERMES_MEET_REALTIME_MODEL"] = realtime_model - if realtime_voice: - env["HERMES_MEET_REALTIME_VOICE"] = realtime_voice - if realtime_instructions: - env["HERMES_MEET_REALTIME_INSTRUCTIONS"] = realtime_instructions - # Resolve the realtime key at SPAWN time, in the parent, where the - # profile secret scope (a contextvar) is still installed. The detached - # child inherits the process environment — NOT the scope — so under a - # multiplexed gateway an in-child os.environ read would see another - # profile's OPENAI_API_KEY (or nothing). Pass it explicitly instead; - # meet_bot checks HERMES_MEET_REALTIME_KEY before OPENAI_API_KEY. - if not realtime_api_key: try: - from agent.secret_scope import get_secret - - realtime_api_key = ( - get_secret("HERMES_MEET_REALTIME_KEY") - or get_secret("OPENAI_API_KEY") - ) - except ImportError: # pragma: no cover — secret_scope is in-repo + (out / name).unlink() + except OSError: pass + + env = {**os.environ, "HERMES_MEET_URL": url, "HERMES_MEET_OUT_DIR": str(out), + "HERMES_MEET_GUEST_NAME": guest_name} + for value, var in ( + (headed and "1", "HERMES_MEET_HEADED"), + (auth_state, "HERMES_MEET_AUTH_STATE"), + (duration, "HERMES_MEET_DURATION"), + (mode, "HERMES_MEET_MODE"), # bot defaults to transcribe when unset (v1 behavior) + (realtime_model, "HERMES_MEET_REALTIME_MODEL"), + (realtime_voice, "HERMES_MEET_REALTIME_VOICE"), + (realtime_instructions, "HERMES_MEET_REALTIME_INSTRUCTIONS"), + ): + if value: + env[var] = value + # Resolve the realtime key at SPAWN time, in the parent, where the profile + # secret scope (a contextvar) is installed. The detached child inherits the + # environment, not the scope — under a multiplexed gateway an in-child + # os.environ read could see another profile's key (or nothing). + if not realtime_api_key: + from agent.secret_scope import get_secret + + realtime_api_key = get_secret("HERMES_MEET_REALTIME_KEY") or get_secret("OPENAI_API_KEY") if realtime_api_key: env["HERMES_MEET_REALTIME_KEY"] = realtime_api_key log_path = out / "bot.log" # Detach: stdin=devnull, stdout/stderr → log file, new session so parent - # signals don't propagate. - log_fh = open(log_path, "ab", buffering=0) - try: + # signals don't propagate. The child owns the log fd after Popen. + with open(log_path, "ab", buffering=0) as log_fh: proc = subprocess.Popen( [sys.executable, "-m", "plugins.google_meet.meet_bot"], - stdin=subprocess.DEVNULL, - stdout=log_fh, - stderr=subprocess.STDOUT, - env=env, - start_new_session=True, - close_fds=True, + stdin=subprocess.DEVNULL, stdout=log_fh, stderr=subprocess.STDOUT, + env=env, start_new_session=True, close_fds=True, ) - finally: - # The subprocess now owns the log fd; we can close ours. - log_fh.close() record = { - "pid": proc.pid, - "meeting_id": meeting_id, - "out_dir": str(out), - "url": url, - "started_at": time.time(), - "session_id": session_id, - "log_path": str(log_path), - "mode": mode, + "pid": proc.pid, "meeting_id": meeting_id, "out_dir": str(out), "url": url, + "started_at": time.time(), "session_id": session_id, "log_path": str(log_path), "mode": mode, } _write_active(record) return {"ok": True, **record} @@ -207,65 +146,42 @@ def status() -> Dict[str, Any]: """Return the current meeting state, or ``{"ok": False, "reason": ...}``.""" active = _read_active() if not active: - return {"ok": False, "reason": "no active meeting"} - + return dict(_NO_ACTIVE) pid = int(active.get("pid", 0)) - alive = _pid_alive(pid) if pid else False - - status_path = Path(active.get("out_dir", "")) / "status.json" - bot_status: Dict[str, Any] = {} - if status_path.is_file(): - try: - bot_status = json.loads(status_path.read_text(encoding="utf-8")) - except Exception: - pass - return { "ok": True, - "alive": alive, + "alive": _pid_alive(pid), "pid": pid, "meetingId": active.get("meeting_id"), "url": active.get("url"), "startedAt": active.get("started_at"), "outDir": active.get("out_dir"), - **bot_status, + **(read_json(Path(active.get("out_dir", "")) / "status.json") or {}), } def transcript(last: Optional[int] = None) -> Dict[str, Any]: - """Read the current transcript file. Returns ok=False if none exists.""" + """Read the current transcript file (empty result if the bot hasn't written one yet).""" active = _read_active() if not active: - return {"ok": False, "reason": "no active meeting"} + return dict(_NO_ACTIVE) tp = Path(active.get("out_dir", "")) / "transcript.txt" - if not tp.is_file(): - return { - "ok": True, - "meetingId": active.get("meeting_id"), - "lines": [], - "total": 0, - "path": str(tp), - } - text = tp.read_text(encoding="utf-8", errors="replace") + text = tp.read_text(encoding="utf-8", errors="replace") if tp.is_file() else "" all_lines = [ln for ln in text.splitlines() if ln.strip()] - lines = all_lines[-last:] if last else all_lines return { "ok": True, "meetingId": active.get("meeting_id"), - "lines": lines, + "lines": all_lines[-last:] if last else all_lines, "total": len(all_lines), "path": str(tp), } def enqueue_say(text: str) -> Dict[str, Any]: - """Append a ``say`` request to the active bot's JSONL queue. + """Append a ``say`` request to ``/say_queue.jsonl`` for the bot's speaker thread. - Returns ``{"ok": False, "reason": ...}`` when no meeting is active or - the active bot is in transcribe-only mode. Otherwise writes a line to - ``/say_queue.jsonl`` that the bot's realtime speaker thread - will consume. + Refused when no meeting is active or the active bot is in transcribe-only mode. """ import uuid @@ -275,15 +191,10 @@ def enqueue_say(text: str) -> Dict[str, Any]: active = _read_active() if not active: - return {"ok": False, "reason": "no active meeting"} + return dict(_NO_ACTIVE) if active.get("mode") != "realtime": - return { - "ok": False, - "reason": ( - "active meeting is in transcribe mode — pass mode='realtime' " - "to meet_join to enable agent speech" - ), - } + return {"ok": False, "reason": "active meeting is in transcribe mode — pass mode='realtime' " + "to meet_join to enable agent speech"} out_dir = Path(active.get("out_dir", "")) if not out_dir.is_dir(): @@ -293,47 +204,35 @@ def enqueue_say(text: str) -> Dict[str, Any]: entry = {"id": uuid.uuid4().hex[:12], "text": text} with queue_path.open("a", encoding="utf-8") as f: f.write(json.dumps(entry) + "\n") - return { - "ok": True, - "meetingId": active.get("meeting_id"), - "enqueued_id": entry["id"], - "queue_path": str(queue_path), - } + return {"ok": True, "meetingId": active.get("meeting_id"), "enqueued_id": entry["id"], + "queue_path": str(queue_path)} def stop(*, reason: str = "requested") -> Dict[str, Any]: - """Signal the active bot to leave cleanly, then clear the active pointer. - - Sends SIGTERM and waits up to 10s for the bot to exit. Falls back to - SIGKILL if the bot doesn't respond. - """ + """SIGTERM the active bot (SIGKILL after 10s), then clear the active pointer.""" active = _read_active() if not active: - return {"ok": False, "reason": "no active meeting"} + return dict(_NO_ACTIVE) pid = int(active.get("pid", 0)) out_dir = active.get("out_dir") - transcript_path = Path(out_dir) / "transcript.txt" if out_dir else None - if pid and _pid_alive(pid): - try: - os.kill(pid, signal.SIGTERM) - except ProcessLookupError: - pass + if _pid_alive(pid): + _kill(pid, signal.SIGTERM) for _ in range(20): if not _pid_alive(pid): break time.sleep(0.5) if _pid_alive(pid): - try: - os.kill(pid, signal.SIGKILL) # windows-footgun: ok — POSIX-only plugin (google_meet registers no-op on Windows; see __init__.py) - except ProcessLookupError: - pass + _kill(pid, signal.SIGKILL) # windows-footgun: ok — POSIX-only plugin (google_meet registers no-op on Windows; see __init__.py) - _clear_active() + try: + (_root() / ".active.json").unlink() + except FileNotFoundError: + pass return { "ok": True, "reason": reason, "meetingId": active.get("meeting_id"), - "transcriptPath": str(transcript_path) if transcript_path else None, + "transcriptPath": str(Path(out_dir) / "transcript.txt") if out_dir else None, } diff --git a/plugins/google_meet/realtime/openai_client.py b/plugins/google_meet/realtime/openai_client.py index 2e84c89125..fac4e2b9b8 100644 --- a/plugins/google_meet/realtime/openai_client.py +++ b/plugins/google_meet/realtime/openai_client.py @@ -1,20 +1,16 @@ """OpenAI Realtime API WebSocket client + file-queue speaker. -This module is the "output" side of the v2 voice bridge: it takes text, -sends it to the OpenAI Realtime API, receives audio deltas back, and -appends the PCM bytes to a file. A separate consumer (the audio -bridge) streams that file into Chrome's fake microphone. - -Designed for simplicity: a single synchronous WebSocket connection per -speaker, per session. The ``websockets`` package is imported lazily so -that importing this module never fails just because the optional dep -is missing. +Output side of the v2 voice bridge: text → OpenAI Realtime → audio deltas +appended as PCM to a file that the audio bridge streams into Chrome's fake +mic. One synchronous WebSocket per speaker/session; ``websockets`` is +imported lazily so importing this module never fails without the optional dep. """ from __future__ import annotations import base64 import json +import threading import time import uuid from pathlib import Path @@ -23,30 +19,21 @@ from typing import Any, Callable, Optional REALTIME_URL = "wss://api.openai.com/v1/realtime" +_TERMINAL_FRAMES = {"response.done", "response.completed", "response.cancelled"} -def _require_websockets(): - """Import ``websockets.sync.client.connect`` or raise with hint.""" + +def _decode_audio(b64: str) -> bytes: try: - from websockets.sync.client import connect as _connect # type: ignore - except ImportError as exc: # pragma: no cover - exercised via test - raise RuntimeError( - "websockets package is required for OpenAI Realtime; " - "install with: pip install websockets" - ) from exc - return _connect + return base64.b64decode(b64) if b64 else b"" + except (ValueError, TypeError): + return b"" class RealtimeSession: """Minimal sync client for the OpenAI Realtime WebSocket API. - Usage: - sess = RealtimeSession(api_key=..., audio_sink_path=Path("out.pcm")) - sess.connect() - sess.speak("Hello team.") - sess.close() - - Thread safety: ``speak`` and ``cancel_response`` may be called from - different threads; a lock serializes WebSocket writes. + ``speak`` and ``cancel_response`` may be called from different threads; + a lock serializes WebSocket writes. """ def __init__( @@ -58,7 +45,6 @@ class RealtimeSession: audio_sink_path: Optional[Path] = None, sample_rate: int = 24000, ) -> None: - import threading as _threading self.api_key = api_key self.model = model self.voice = voice @@ -66,42 +52,38 @@ class RealtimeSession: self.audio_sink_path = Path(audio_sink_path) if audio_sink_path else None self.sample_rate = sample_rate self._ws: Any = None - self._send_lock = _threading.Lock() - self._last_response_id: Optional[str] = None + self._send_lock = threading.Lock() # Public counters for status reporting. self.audio_bytes_out: int = 0 self.last_audio_out_at: Optional[float] = None - # ── lifecycle ───────────────────────────────────────────────────────── - def connect(self) -> None: - """Open WS and send session.update with voice+instructions.""" - connect = _require_websockets() + """Open the WS and send ``session.update`` with voice + instructions.""" + try: + from websockets.sync.client import connect # type: ignore + except ImportError as exc: # pragma: no cover - exercised via test + raise RuntimeError( + "websockets package is required for OpenAI Realtime; " + "install with: pip install websockets" + ) from exc url = f"{REALTIME_URL}?model={self.model}" - headers = [ - ("Authorization", f"Bearer {self.api_key}"), - ("OpenAI-Beta", "realtime=v1"), - ] - # websockets.sync.client.connect accepts either additional_headers= - # (newer) or extra_headers= depending on version; try the newer - # name first and fall back. + headers = [("Authorization", f"Bearer {self.api_key}"), ("OpenAI-Beta", "realtime=v1")] + # Newer websockets takes additional_headers=, older extra_headers=. try: self._ws = connect(url, additional_headers=headers) except TypeError: self._ws = connect(url, extra_headers=headers) - self._send_json( - { - "type": "session.update", - "session": { - "voice": self.voice, - "instructions": self.instructions, - "modalities": ["audio", "text"], - "output_audio_format": "pcm16", - "input_audio_format": "pcm16", - }, - } - ) + self._send_json({ + "type": "session.update", + "session": { + "voice": self.voice, + "instructions": self.instructions, + "modalities": ["audio", "text"], + "output_audio_format": "pcm16", + "input_audio_format": "pcm16", + }, + }) def close(self) -> None: if self._ws is not None: @@ -111,107 +93,54 @@ class RealtimeSession: pass self._ws = None - # ── speaking ────────────────────────────────────────────────────────── - def speak(self, text: str, timeout: float = 30.0) -> dict: - """Send ``text`` and accumulate the audio response. + """Send ``text`` and append the audio response to ``audio_sink_path``. - Audio deltas are base64-decoded and appended to - ``audio_sink_path`` (opened 'ab' and closed per call, so a - separate streaming reader can consume whatever is there). + The sink is opened 'ab' and closed per call so a streaming reader can + consume whatever is there. Frames other than audio deltas, terminal + response events and errors are ignored. """ if self._ws is None: raise RuntimeError("RealtimeSession.connect() must be called first") start = time.monotonic() - - self._send_json( - { - "type": "conversation.item.create", - "item": { - "type": "message", - "role": "user", - "content": [{"type": "input_text", "text": text}], - }, - } - ) - self._send_json( - { - "type": "response.create", - "response": {"modalities": ["audio"]}, - } - ) + self._send_json({ + "type": "conversation.item.create", + "item": {"type": "message", "role": "user", "content": [{"type": "input_text", "text": text}]}, + }) + self._send_json({"type": "response.create", "response": {"modalities": ["audio"]}}) bytes_written = 0 sink_fp = None if self.audio_sink_path is not None: self.audio_sink_path.parent.mkdir(parents=True, exist_ok=True) sink_fp = open(self.audio_sink_path, "ab") - try: while True: - remaining = timeout - (time.monotonic() - start) - if remaining <= 0: - raise TimeoutError( - f"realtime response did not complete within {timeout}s" - ) - raw = self._recv(timeout=remaining) - if raw is None: - # Connection closed by peer. + frame = self._recv_frame(start + timeout, timeout) + if frame is None: # connection closed by peer break - try: - frame = json.loads(raw) if isinstance(raw, (str, bytes, bytearray)) else raw - except (TypeError, ValueError): - continue - if not isinstance(frame, dict): - continue ftype = frame.get("type") - if ftype == "response.audio.delta": - b64 = frame.get("delta") or frame.get("audio") or "" - if b64 and sink_fp is not None: - try: - chunk = base64.b64decode(b64) - except (ValueError, TypeError): - chunk = b"" - if chunk: - sink_fp.write(chunk) - sink_fp.flush() - bytes_written += len(chunk) - self.audio_bytes_out += len(chunk) - self.last_audio_out_at = time.time() - elif ftype == "response.created": - rid = (frame.get("response") or {}).get("id") - if rid: - self._last_response_id = rid - elif ftype in {"response.done", "response.completed", "response.cancelled"}: + if ftype in _TERMINAL_FRAMES: break - elif ftype == "error": - err = frame.get("error") or frame - raise RuntimeError(f"realtime error: {err}") - # All other frames (response.created, response.output_item.*, - # response.audio_transcript.delta, rate_limits.updated, ...) - # are ignored for v2. + if ftype == "error": + raise RuntimeError(f"realtime error: {frame.get('error') or frame}") + if ftype == "response.audio.delta" and sink_fp is not None: + chunk = _decode_audio(frame.get("delta") or frame.get("audio") or "") + if chunk: + sink_fp.write(chunk) + sink_fp.flush() + bytes_written += len(chunk) + self.audio_bytes_out += len(chunk) + self.last_audio_out_at = time.time() finally: if sink_fp is not None: sink_fp.close() - duration_ms = (time.monotonic() - start) * 1000.0 - return { - "ok": True, - "bytes_written": bytes_written, - "duration_ms": duration_ms, - } - - # ── ws plumbing ─────────────────────────────────────────────────────── + return {"ok": True, "bytes_written": bytes_written, "duration_ms": (time.monotonic() - start) * 1000.0} def cancel_response(self) -> bool: - """Interrupt the in-flight response (barge-in). - - Sends ``response.cancel`` on the current WebSocket so the model - stops generating audio immediately. Safe to call at any time; - returns True if a cancel was actually sent, False when there's - nothing to cancel or the socket isn't open. - """ + """Barge-in: send ``response.cancel``. True if sent, False if nothing to cancel / socket closed.""" if self._ws is None: return False try: @@ -225,25 +154,35 @@ class RealtimeSession: with self._send_lock: self._ws.send(json.dumps(payload)) - def _recv(self, timeout: Optional[float] = None): + def _recv_frame(self, deadline: float, timeout: float) -> Optional[dict]: + """Next dict frame before *deadline* (monotonic), ``None`` once the peer closes. + + Non-dict / unparseable frames are skipped; TimeoutError past the deadline. + """ assert self._ws is not None - try: - if timeout is None: - return self._ws.recv() - return self._ws.recv(timeout=timeout) - except TypeError: - # Older websockets may not accept timeout kwarg. - return self._ws.recv() + while True: + remaining = deadline - time.monotonic() + if remaining <= 0: + raise TimeoutError(f"realtime response did not complete within {timeout}s") + try: + raw = self._ws.recv(timeout=remaining) + except TypeError: # older websockets: no timeout kwarg + raw = self._ws.recv() + if raw is None: + return None + try: + frame = json.loads(raw) if isinstance(raw, (str, bytes, bytearray)) else raw + except (TypeError, ValueError): + continue + if isinstance(frame, dict): + return frame class RealtimeSpeaker: """File-based JSONL queue wrapper around :class:`RealtimeSession`. - Each line in ``queue_path`` is a JSON object of the form - ``{"id": "", "text": "..."}``. Processed lines are appended - to ``processed_path`` (if set) and then removed from the queue; - if ``processed_path`` is ``None``, processed lines are simply - dropped. + Each queue line is ``{"id": "", "text": "..."}``. Processed lines + are appended to ``processed_path`` (if set) and removed from the queue. """ def __init__( @@ -256,36 +195,26 @@ class RealtimeSpeaker: self.queue_path = Path(queue_path) self.processed_path = Path(processed_path) if processed_path else None - # ── helpers ────────────────────────────────────────────────────────── - def _read_queue(self) -> list[dict]: + """Parse the JSONL queue, skipping blank/malformed lines; entries lacking an ``id`` get one.""" if not self.queue_path.exists(): return [] out: list[dict] = [] for line in self.queue_path.read_text(encoding="utf-8").splitlines(): - line = line.strip() - if not line: - continue try: - entry = json.loads(line) + entry = json.loads(line) if line.strip() else None except ValueError: continue - if not isinstance(entry, dict): - continue - if "id" not in entry: - entry["id"] = str(uuid.uuid4()) - out.append(entry) + if isinstance(entry, dict): + entry.setdefault("id", str(uuid.uuid4())) + out.append(entry) return out def _rewrite_queue(self, remaining: list[dict]) -> None: - if not remaining: - # Keep the file but empty — consumers may be watching for - # new writes via mtime, and delete-then-recreate is a race. - self.queue_path.write_text("", encoding="utf-8") - return - self.queue_path.write_text( - "\n".join(json.dumps(e) for e in remaining) + "\n", encoding="utf-8" - ) + # Always keep the file (empty when drained): consumers may watch its + # mtime, and delete-then-recreate is a race. + body = "".join(json.dumps(e) + "\n" for e in remaining) + self.queue_path.write_text(body, encoding="utf-8") def _append_processed(self, entry: dict, result: dict) -> None: if self.processed_path is None: @@ -295,20 +224,13 @@ class RealtimeSpeaker: with open(self.processed_path, "a", encoding="utf-8") as fp: fp.write(json.dumps(record) + "\n") - # ── main loop ──────────────────────────────────────────────────────── - - def run_until_stopped( - self, - stop_fn: Callable[[], bool], - poll_interval: float = 0.5, - ) -> None: + def run_until_stopped(self, stop_fn: Callable[[], bool], poll_interval: float = 0.5) -> None: while not stop_fn(): entries = self._read_queue() if not entries: time.sleep(poll_interval) continue - # Process one at a time; re-check the queue file after each - # speak() call because new entries may have arrived. + # One entry per iteration: the queue may grow while we speak. head = entries[0] text = (head.get("text") or "").strip() if text: @@ -320,13 +242,10 @@ class RealtimeSpeaker: result = {"ok": True, "bytes_written": 0, "duration_ms": 0.0} self._append_processed(head, result) - # Re-read the queue from disk in case it was appended to - # while we were speaking, then drop the head. + # Re-read from disk (new entries may have arrived), then drop the + # head — by position when it's still first, else by id. latest = self._read_queue() if latest and latest[0].get("id") == head.get("id"): self._rewrite_queue(latest[1:]) else: - # Fallback: drop-by-id anywhere in the queue. - self._rewrite_queue( - [e for e in latest if e.get("id") != head.get("id")] - ) + self._rewrite_queue([e for e in latest if e.get("id") != head.get("id")]) diff --git a/plugins/google_meet/tools.py b/plugins/google_meet/tools.py index 034116b88a..61df319d22 100644 --- a/plugins/google_meet/tools.py +++ b/plugins/google_meet/tools.py @@ -1,14 +1,10 @@ """Agent-facing tools for the google_meet plugin. -Tools: - meet_join — join a Google Meet URL (spawns Playwright bot locally - OR on a remote node host via node=) - meet_status — report bot liveness + transcript progress - meet_transcript — read the current transcript (optional last-N) + meet_join — join a Meet URL (locally, or on a remote node via node=) + meet_status — bot liveness + transcript progress + meet_transcript — read the transcript (optional last-N) meet_leave — signal the bot to leave cleanly - meet_say — (v2) speak text through the realtime audio bridge. - Requires the active meeting to have been joined with - mode='realtime'. + meet_say — speak text through the realtime bridge (mode='realtime' only) """ from __future__ import annotations @@ -19,21 +15,11 @@ from typing import Any, Dict, Optional from plugins.google_meet import process_manager as pm -# --------------------------------------------------------------------------- -# Runtime gate -# --------------------------------------------------------------------------- - def check_meet_requirements() -> bool: - """Return True when the plugin can actually run LOCALLY. + """True when the plugin can run LOCALLY: Linux/macOS + importable ``playwright``. - Gates on: - * Python ``playwright`` package importable - * the plugin being on a supported platform (Linux or macOS) - - Note: remote-node operation (``node=``) only needs the - ``websockets`` dep on the gateway side — Chromium lives on the node. - But the plugin-level gate keeps the v1 semantics; individual tool - handlers relax the requirement when a node is addressed. + Remote-node operation only needs ``websockets`` on the gateway side; the + handlers relax this gate themselves when a node is addressed. """ import platform as _p if _p.system().lower() not in {"linux", "darwin"}: @@ -45,36 +31,23 @@ def check_meet_requirements() -> bool: return True -# --------------------------------------------------------------------------- -# Node client helper -# --------------------------------------------------------------------------- - -def _resolve_node_client(node: Optional[str]): - """Return (NodeClient, node_name) for *node*, or (None, None) to run local. - - Raises RuntimeError with a readable message if the node is named but - unresolvable, so the handler can surface a clear error to the agent. - """ - if node is None or node == "": - return None, None +def resolve_node(node: str): + """``(NodeClient, node_name)`` for *node* (``'auto'`` = the sole registered node), or ``(None, None)``.""" from plugins.google_meet.node.registry import NodeRegistry from plugins.google_meet.node.client import NodeClient - reg = NodeRegistry() - entry = reg.resolve(node if node != "auto" else None) + entry = NodeRegistry().resolve(node if node != "auto" else None) if entry is None: - raise RuntimeError( - f"no registered meet node matches {node!r} — " - "run `hermes meet node approve ` first" - ) - client = NodeClient(url=entry["url"], token=entry["token"]) - return client, entry.get("name") + return None, None + return NodeClient(url=entry["url"], token=entry["token"]), entry.get("name") # --------------------------------------------------------------------------- # Schemas # --------------------------------------------------------------------------- +_NODE_PROP = {"type": "string"} + MEET_JOIN_SCHEMA: Dict[str, Any] = { "name": "meet_join", "description": ( @@ -89,12 +62,7 @@ MEET_JOIN_SCHEMA: Dict[str, Any] = { "parameters": { "type": "object", "properties": { - "url": { - "type": "string", - "description": ( - "Full https://meet.google.com/... URL. Required." - ), - }, + "url": {"type": "string", "description": "Full https://meet.google.com/... URL. Required."}, "mode": { "type": "string", "enum": ["transcribe", "realtime"], @@ -106,10 +74,7 @@ MEET_JOIN_SCHEMA: Dict[str, Any] = { }, "guest_name": { "type": "string", - "description": ( - "Display name to use when joining as guest. Defaults to " - "'Hermes Agent'." - ), + "description": "Display name to use when joining as guest. Defaults to 'Hermes Agent'.", }, "duration": { "type": "string", @@ -120,10 +85,7 @@ MEET_JOIN_SCHEMA: Dict[str, Any] = { }, "headed": { "type": "boolean", - "description": ( - "Run Chromium headed instead of headless (debug only). " - "Default false." - ), + "description": "Run Chromium headed instead of headless (debug only). Default false.", }, "node": { "type": "string", @@ -149,13 +111,7 @@ MEET_STATUS_SCHEMA: Dict[str, Any] = { "has joined, is sitting in the lobby, number of transcript lines " "captured, and last-caption timestamp." ), - "parameters": { - "type": "object", - "properties": { - "node": {"type": "string"}, - }, - "additionalProperties": False, - }, + "parameters": {"type": "object", "properties": {"node": _NODE_PROP}, "additionalProperties": False}, } MEET_TRANSCRIPT_SCHEMA: Dict[str, Any] = { @@ -177,7 +133,7 @@ MEET_TRANSCRIPT_SCHEMA: Dict[str, Any] = { ), "minimum": 1, }, - "node": {"type": "string"}, + "node": _NODE_PROP, }, "additionalProperties": False, }, @@ -190,13 +146,7 @@ MEET_LEAVE_SCHEMA: Dict[str, Any] = { "finalize the transcript file. Safe to call when no meeting is " "active — returns ok=false with a reason." ), - "parameters": { - "type": "object", - "properties": { - "node": {"type": "string"}, - }, - "additionalProperties": False, - }, + "parameters": {"type": "object", "properties": {"node": _NODE_PROP}, "additionalProperties": False}, } MEET_SAY_SCHEMA: Dict[str, Any] = { @@ -211,10 +161,7 @@ MEET_SAY_SCHEMA: Dict[str, Any] = { ), "parameters": { "type": "object", - "properties": { - "text": {"type": "string", "description": "Text to speak."}, - "node": {"type": "string"}, - }, + "properties": {"text": {"type": "string", "description": "Text to speak."}, "node": _NODE_PROP}, "required": ["text"], "additionalProperties": False, }, @@ -233,6 +180,24 @@ def _err(msg: str, **extra) -> str: return _json({"success": False, "error": msg, **extra}) +def _dispatch(node: Optional[str], op: str, remote, local) -> str: + """Run *remote(client)* on the addressed node, else *local()*; wrap as a tool result.""" + if not node: + res = local() + return _json({"success": bool(res.get("ok")), **res}) + client, node_name = resolve_node(node) + if client is None: + return _err( + f"no registered meet node matches {node!r} — " + "run `hermes meet node approve ` first" + ) + try: + res = remote(client) + except Exception as e: + return _err(f"remote node {op} failed: {e}", node=node_name) + return _json({"success": bool(res.get("ok")), "node": node_name, **res}) + + def handle_meet_join(args: Dict[str, Any], **_kw) -> str: url = (args.get("url") or "").strip() if not url: @@ -241,108 +206,55 @@ def handle_meet_join(args: Dict[str, Any], **_kw) -> str: if mode not in {"transcribe", "realtime"}: return _err(f"mode must be 'transcribe' or 'realtime' (got {mode!r})") - node = args.get("node") - try: - client, node_name = _resolve_node_client(node) - except RuntimeError as e: - return _err(str(e)) - - if client is not None: - # Remote path — delegate to the node host. - try: - res = client.start_bot( - url=url, - guest_name=str(args.get("guest_name") or "Hermes Agent"), - duration=str(args.get("duration")) if args.get("duration") else None, - headed=bool(args.get("headed", False)), - mode=mode, - ) - return _json({"success": bool(res.get("ok")), "node": node_name, **res}) - except Exception as e: - return _err(f"remote node start_bot failed: {e}", node=node_name) - - # Local path — same as v1, with v2 params. - if not check_meet_requirements(): - return _err( - "google_meet plugin prerequisites missing — install with " - "`pip install playwright && python -m playwright install " - "chromium`. Plugin is supported on Linux and macOS only." - ) - res = pm.start( + common: Dict[str, Any] = dict( url=url, - headed=bool(args.get("headed", False)), guest_name=str(args.get("guest_name") or "Hermes Agent"), duration=str(args.get("duration")) if args.get("duration") else None, + headed=bool(args.get("headed", False)), mode=mode, ) - return _json({"success": bool(res.get("ok")), **res}) + + def _local(): + if not check_meet_requirements(): + return { + "ok": False, + "error": ( + "google_meet plugin prerequisites missing — install with " + "`pip install playwright && python -m playwright install " + "chromium`. Plugin is supported on Linux and macOS only." + ), + } + return pm.start(**common) + + return _dispatch(args.get("node"), "start_bot", lambda c: c.start_bot(**common), _local) def handle_meet_status(args: Dict[str, Any], **_kw) -> str: - try: - client, node_name = _resolve_node_client(args.get("node")) - except RuntimeError as e: - return _err(str(e)) - if client is not None: - try: - res = client.status() - return _json({"success": bool(res.get("ok")), "node": node_name, **res}) - except Exception as e: - return _err(f"remote node status failed: {e}", node=node_name) - res = pm.status() - return _json({"success": bool(res.get("ok")), **res}) + return _dispatch(args.get("node"), "status", lambda c: c.status(), pm.status) def handle_meet_transcript(args: Dict[str, Any], **_kw) -> str: - last = args.get("last") try: - last_i = int(last) if last is not None else None - if last_i is not None and last_i < 1: - last_i = None + last = int(args["last"]) if args.get("last") is not None else None except (TypeError, ValueError): - last_i = None - try: - client, node_name = _resolve_node_client(args.get("node")) - except RuntimeError as e: - return _err(str(e)) - if client is not None: - try: - res = client.transcript(last=last_i) - return _json({"success": bool(res.get("ok")), "node": node_name, **res}) - except Exception as e: - return _err(f"remote node transcript failed: {e}", node=node_name) - res = pm.transcript(last=last_i) - return _json({"success": bool(res.get("ok")), **res}) + last = None + if last is not None and last < 1: + last = None + return _dispatch( + args.get("node"), "transcript", + lambda c: c.transcript(last=last), lambda: pm.transcript(last=last), + ) def handle_meet_leave(args: Dict[str, Any], **_kw) -> str: - try: - client, node_name = _resolve_node_client(args.get("node")) - except RuntimeError as e: - return _err(str(e)) - if client is not None: - try: - res = client.stop() - return _json({"success": bool(res.get("ok")), "node": node_name, **res}) - except Exception as e: - return _err(f"remote node stop failed: {e}", node=node_name) - res = pm.stop(reason="agent called meet_leave") - return _json({"success": bool(res.get("ok")), **res}) + return _dispatch( + args.get("node"), "stop", + lambda c: c.stop(), lambda: pm.stop(reason="agent called meet_leave"), + ) def handle_meet_say(args: Dict[str, Any], **_kw) -> str: text = (args.get("text") or "").strip() if not text: return _err("text is required") - try: - client, node_name = _resolve_node_client(args.get("node")) - except RuntimeError as e: - return _err(str(e)) - if client is not None: - try: - res = client.say(text) - return _json({"success": bool(res.get("ok")), "node": node_name, **res}) - except Exception as e: - return _err(f"remote node say failed: {e}", node=node_name) - res = pm.enqueue_say(text) - return _json({"success": bool(res.get("ok")), **res}) + return _dispatch(args.get("node"), "say", lambda c: c.say(text), lambda: pm.enqueue_say(text)) diff --git a/tests/plugins/test_google_meet_audio.py b/tests/plugins/test_google_meet_audio.py index 3dbce7372e..3703191f86 100644 --- a/tests/plugins/test_google_meet_audio.py +++ b/tests/plugins/test_google_meet_audio.py @@ -138,22 +138,6 @@ _BH_ABSENT = ( # --------------------------------------------------------------------------- -# --------------------------------------------------------------------------- -# chrome_fake_audio_flags -# --------------------------------------------------------------------------- - - -def test_chrome_fake_audio_flags_linux(): - from plugins.google_meet.audio_bridge import chrome_fake_audio_flags - - with patch("plugins.google_meet.audio_bridge.platform.system", - return_value="Linux"): - flags = chrome_fake_audio_flags( - {"platform": "linux", "device_name": "hermes_meet_src"} - ) - assert "--use-fake-ui-for-media-stream" in flags - - def test_property_access_before_setup_raises(): from plugins.google_meet.audio_bridge import AudioBridge diff --git a/tests/plugins/test_google_meet_plugin.py b/tests/plugins/test_google_meet_plugin.py index 74b8665193..d12ea77d43 100644 --- a/tests/plugins/test_google_meet_plugin.py +++ b/tests/plugins/test_google_meet_plugin.py @@ -262,12 +262,12 @@ def test_looks_like_human_speaker(): def test_detect_admission_returns_false_on_error(): - from plugins.google_meet.meet_bot import _detect_admission + from plugins.google_meet.meet_bot import _ADMISSION_PROBE_JS, _probe class _FakePage: def evaluate(self, _js): raise RuntimeError("boom") - assert _detect_admission(_FakePage()) is False + assert _probe(_FakePage(), _ADMISSION_PROBE_JS) is False # ---------------------------------------------------------------------------