diff --git a/plugins/context_engine/__init__.py b/plugins/context_engine/__init__.py index 12e2f4b7c3..f3e7450898 100644 --- a/plugins/context_engine/__init__.py +++ b/plugins/context_engine/__init__.py @@ -1,10 +1,6 @@ -"""Context engine plugin discovery. - -Scans ``plugins/context_engine//`` for ``ContextEngine`` implementations. Engines live -in the repo (always available, no user install) and are separate from the general plugin -system. Only one is active (``context.engine`` in config.yaml; default ``"compressor"``, -the built-in ContextCompressor). -""" +"""Context engine plugin discovery: ``plugins/context_engine//`` → ``ContextEngine``. +Engines ship in the repo, separate from the general plugin system; only one is active +(``context.engine`` in config.yaml; default ``"compressor"``, the built-in ContextCompressor).""" from __future__ import annotations @@ -21,12 +17,9 @@ _CONTEXT_ENGINE_PLUGINS_DIR = Path(__file__).parent def discover_context_engines() -> List[Tuple[str, str, bool]]: """Return ``[(name, description, is_available), ...]`` for every bundled engine.""" - return [ - ( - child.name, - _loader.read_plugin_description(child), - _loader.probe_availability(lambda c=child: _load_engine_from_dir(c))) - for child in _loader.iter_plugin_dirs(_CONTEXT_ENGINE_PLUGINS_DIR)] + return [(child.name, _loader.read_plugin_description(child), + _loader.probe_availability(lambda c=child: _load_engine_from_dir(c))) + for child in _loader.iter_plugin_dirs(_CONTEXT_ENGINE_PLUGINS_DIR)] def load_context_engine(name: str) -> Optional["ContextEngine"]: # noqa: F821 @@ -56,8 +49,8 @@ def _load_engine_from_dir(engine_dir: Path) -> Optional["ContextEngine"]: # noq class _EngineCollector(_loader.NoopPluginContext): - """Fake plugin context capturing register_context_engine; forwards register_command - to the global plugin command registry so engine slash commands behave like plugin ones.""" + """Captures register_context_engine; forwards register_command to the global plugin command + registry so engine slash commands behave like plugin ones.""" def __init__(self, engine_name: str = ""): self.engine = None @@ -70,9 +63,8 @@ class _EngineCollector(_loader.NoopPluginContext): def register_command(self, name: str, handler, description: str = "", args_hint: str = "") -> None: clean = (name or "").lower().strip().lstrip("/").replace(" ", "-") if not clean: - logger.warning( - "Context engine '%s' tried to register a command with an empty name.", - self._engine_name) + logger.warning("Context engine '%s' tried to register a command with an empty name.", + self._engine_name) return conflict = "Context engine '%s' tried to register command '/%s' which %s Skipping." @@ -91,12 +83,9 @@ class _EngineCollector(_loader.NoopPluginContext): logger.warning(conflict, self._engine_name, clean, "is already registered by a plugin.") return manager._plugin_commands[clean] = { - "handler": handler, - "description": description or "Context engine command", - "plugin": f"context-engine:{self._engine_name}", - "args_hint": (args_hint or "").strip()} + "handler": handler, "description": description or "Context engine command", + "plugin": f"context-engine:{self._engine_name}", "args_hint": (args_hint or "").strip()} self._registered_commands.append(clean) logger.debug("Context engine '%s' registered command: /%s", self._engine_name, clean) except Exception as exc: - logger.debug( - "Context engine '%s' could not register /%s: %s", self._engine_name, clean, exc) + logger.debug("Context engine '%s' could not register /%s: %s", self._engine_name, clean, exc) diff --git a/plugins/cron_providers/__init__.py b/plugins/cron_providers/__init__.py index 799ebfd38e..ecd7eebf52 100644 --- a/plugins/cron_providers/__init__.py +++ b/plugins/cron_providers/__init__.py @@ -1,11 +1,7 @@ -"""Cron scheduler provider plugin discovery. - -Scans bundled ``plugins/cron_providers//`` then user ``$HERMES_HOME/plugins//`` -for ``CronScheduler`` implementations (bundled wins on name collision). The built-in -``InProcessCronScheduler`` is core (``cron/scheduler_provider.py``), not discovered here, -so the fallback can never be accidentally removed; only one provider is active -(``cron.provider`` in config.yaml, empty = built-in). -""" +"""Cron scheduler provider discovery: bundled ``plugins/cron_providers//`` then user +``$HERMES_HOME/plugins//`` (bundled wins on collision). The built-in InProcessCronScheduler +is core, not discovered here, so the fallback can't be removed; one provider is active +(``cron.provider`` in config.yaml, empty = built-in).""" from __future__ import annotations @@ -22,7 +18,6 @@ _CRON_PLUGINS_DIR = Path(__file__).parent _USER_NAMESPACE = "_hermes_user_cron" - def _is_cron_provider_dir(path: Path) -> bool: """Cheap text heuristic: ``__init__.py`` mentions the cron scheduler contract.""" init_file = path / "__init__.py" @@ -40,10 +35,8 @@ def _user_provider_dirs() -> List[Path]: user_dir = _loader.user_plugins_dir() if not user_dir: return [] - return [ - child for child in sorted(user_dir.iterdir()) - if child.is_dir() and not child.name.startswith(("_", ".")) and _is_cron_provider_dir(child) - ] + return [child for child in sorted(user_dir.iterdir()) + if child.is_dir() and not child.name.startswith(("_", ".")) and _is_cron_provider_dir(child)] def _iter_provider_dirs() -> List[Tuple[str, Path]]: @@ -67,12 +60,9 @@ def find_provider_dir(name: str) -> Optional[Path]: def discover_cron_schedulers() -> List[Tuple[str, str, bool]]: """Return ``[(name, description, is_available), ...]`` for all discovered providers.""" - return [ - ( - name, - _loader.read_plugin_description(child), - _loader.probe_availability(lambda c=child: _load_provider_from_dir(c))) - for name, child in _iter_provider_dirs()] + return [(name, _loader.read_plugin_description(child), + _loader.probe_availability(lambda c=child: _load_provider_from_dir(c))) + for name, child in _iter_provider_dirs()] def load_cron_scheduler(name: str) -> Optional["CronScheduler"]: # noqa: F821 @@ -94,9 +84,7 @@ def _load_provider_from_dir(provider_dir: Path) -> Optional["CronScheduler"]: # is_bundled = _CRON_PLUGINS_DIR in provider_dir.parents or provider_dir.parent == _CRON_PLUGINS_DIR module_name = f"plugins.cron_providers.{name}" if is_bundled else f"{_USER_NAMESPACE}.{name}" mod = _loader.load_plugin_module( - module_name, provider_dir, - parents=("plugins", "plugins.cron_providers"), - logger=logger, + module_name, provider_dir, parents=("plugins", "plugins.cron_providers"), logger=logger, synthetic_namespace=None if is_bundled else _USER_NAMESPACE) if mod is None: return None diff --git a/plugins/cron_providers/chronos/__init__.py b/plugins/cron_providers/chronos/__init__.py index b5ce44bdd2..764540ac28 100644 --- a/plugins/cron_providers/chronos/__init__.py +++ b/plugins/cron_providers/chronos/__init__.py @@ -1,11 +1,9 @@ """Chronos — NAS-mediated managed cron provider (scale-to-zero). -Instead of a 60s in-process ticker, asks NAS to arm one external one-shot per job at its -next-fire time; NAS calls back ``/api/cron/fire`` and the job re-arms after running. -Between fires the machine is truly idle: start() never blocks or spawns a periodic wake, -and reconcile only runs on a warm process (start / on_jobs_changed / piggybacked on a fire). -Chronos names no scheduler vendor and holds no scheduler credentials — it speaks only to -NAS's ``agent-cron`` endpoints with the agent's existing Nous token. +Instead of a 60s ticker, asks NAS to arm one external one-shot per job at its next-fire time; +NAS calls back ``/api/cron/fire`` and the job re-arms after running. start() never blocks or +spawns a periodic wake; reconcile runs only on a warm process (start / on_jobs_changed / fire). +Holds no scheduler credentials — speaks only to NAS ``agent-cron`` endpoints with the Nous token. Inert unless ``cron.provider: chronos``. Wire contract: ``docs/chronos-managed-cron-contract.md``. """ @@ -33,8 +31,7 @@ class ChronosCronScheduler(CronScheduler): """NAS-mediated external cron provider.""" def __init__(self) -> None: - # Best-effort cache of job_id -> fire_at armed via NAS; reconcile rebuilds - # desired state from jobs.json, so a cold process simply re-arms (idempotent). + # Best-effort job_id -> fire_at cache; a cold process simply re-arms (idempotent). self._armed: Dict[str, str] = {} self._lock = threading.Lock() self._client = None # lazily constructed (no network in is_available) @@ -44,12 +41,11 @@ class ChronosCronScheduler(CronScheduler): return "chronos" def is_available(self) -> bool: - """Config presence only — NO network. Needs portal URL, a publicly reachable - callback URL (for NAS->agent fires) and a Nous login; otherwise the resolver - falls back to the built-in ticker.""" + """Config presence only — NO network: portal URL, publicly reachable callback URL and a + Nous login; otherwise the resolver falls back to the built-in ticker.""" if not (_cfg("cron", "chronos", "portal_url") and _cfg("cron", "chronos", "callback_url")): return False - # Stored-token presence only (no refresh); the refresh-aware token is resolved at provision time. + # Stored-token presence only (no refresh); refresh-aware token resolved at provision time. try: from hermes_cli.auth import get_provider_auth_state return bool((get_provider_auth_state("nous") or {}).get("access_token")) @@ -62,12 +58,9 @@ class ChronosCronScheduler(CronScheduler): self._client = NasCronClient(_cfg("cron", "chronos", "portal_url")) return self._client - # -- lifecycle -------------------------------------------------------- - def start(self, stop_event, *, adapters=None, loop=None, interval=60): """Arm all enabled jobs via NAS, then RETURN — no loop, no periodic wake (scale-to-zero).""" - # A new lifecycle cannot prove what an interrupted prior process did: - # classify those attempts unknown for audit only, never requeue here. + # A new lifecycle can't prove what an interrupted process did: classify unknown, never requeue. self.recover_interrupted() try: self.reconcile() @@ -87,20 +80,16 @@ class ChronosCronScheduler(CronScheduler): """Arm the first one-shot for a new job; may raise so creation can report it.""" self._arm_one_shot(job) - # -- arming ----------------------------------------------------------- - def _arm_one_shot(self, job: Dict[str, Any]) -> None: - """Arm exactly one one-shot at next_run_at. The agent computes the time; NAS is the - dumb executor. dedup_key=(job_id, fire_at) makes re-arming the same fire a no-op.""" + """Arm one one-shot at next_run_at (agent computes the time; NAS executes). + dedup_key=(job_id, fire_at) makes re-arming the same fire a no-op.""" job_id = job["id"] fire_at = job.get("next_run_at") if not fire_at: return self._get_client().provision( - job_id=job_id, - fire_at=fire_at, - agent_callback_url=str(_cfg("cron", "chronos", "callback_url") or ""), - dedup_key=f"{job_id}:{fire_at}") + job_id=job_id, fire_at=fire_at, dedup_key=f"{job_id}:{fire_at}", + agent_callback_url=str(_cfg("cron", "chronos", "callback_url") or "")) with self._lock: self._armed[job_id] = fire_at @@ -119,16 +108,14 @@ class ChronosCronScheduler(CronScheduler): self._armed.pop(job_id, None) def _list_armed(self) -> Dict[str, str]: - """Observed armed one-shots (job_id -> fire_at): in-memory map when warm, else ask - NAS. If the NAS list fails return {} — reconcile then re-arms idempotently.""" + """Armed one-shots (job_id -> fire_at): in-memory map when warm, else ask NAS ({} on + failure — reconcile then re-arms idempotently).""" with self._lock: if self._armed: return dict(self._armed) try: - observed = { - item["job_id"]: item.get("fire_at", "") - for item in self._get_client().list_armed() - if item.get("job_id")} + observed = {item["job_id"]: item.get("fire_at", "") + for item in self._get_client().list_armed() if item.get("job_id")} with self._lock: self._armed.update(observed) return observed @@ -141,8 +128,7 @@ class ChronosCronScheduler(CronScheduler): from cron.jobs import get_job, load_jobs desired: Dict[str, str] = { - j["id"]: j["next_run_at"] - for j in load_jobs() + j["id"]: j["next_run_at"] for j in load_jobs() if j.get("enabled") and j.get("next_run_at") and j.get("state") != "paused"} observed = self._list_armed() @@ -152,22 +138,19 @@ class ChronosCronScheduler(CronScheduler): if job: self._arm_logged(job, f"arm job {job_id}") - for job_id in list(observed.keys()): - if job_id not in desired: - try: - self._cancel(job_id) - except Exception as e: - logger.warning("Chronos failed to cancel orphan %s: %s", job_id, e) + for job_id in observed.keys() - desired.keys(): + try: + self._cancel(job_id) + except Exception as e: + logger.warning("Chronos failed to cancel orphan %s: %s", job_id, e) - # No ``fire_due`` override on purpose: ``provider_supports_split_fire`` treats ANY - # override as the legacy single-phase signal, which would opt Chronos out of claim - # admission, duplicate detection, and cancel-aware drain on the fire webhook. + # No ``fire_due`` override on purpose: ``provider_supports_split_fire`` treats ANY override as + # the legacy single-phase signal, opting Chronos out of claim admission and cancel-aware drain. def fire_claimed( self, claimed_job: dict, *, adapters: Any = None, loop: Any = None, cancel_event: Any = None ) -> bool: - ran = super().fire_claimed( - claimed_job, adapters=adapters, loop=loop, cancel_event=cancel_event) + ran = super().fire_claimed(claimed_job, adapters=adapters, loop=loop, cancel_event=cancel_event) if ran: from cron.jobs import get_job job_id = claimed_job["id"] diff --git a/plugins/cron_providers/chronos/_nas_client.py b/plugins/cron_providers/chronos/_nas_client.py index 4d6fa34a2d..8048065e87 100644 --- a/plugins/cron_providers/chronos/_nas_client.py +++ b/plugins/cron_providers/chronos/_nas_client.py @@ -1,9 +1,6 @@ -"""Thin HTTP client for the agent -> NAS ``agent-cron`` endpoints (Chronos). - -The agent only asks NAS to "arm a one-shot at time T" / "cancel" / "list", authenticated -with its existing Nous Portal access token (no new secret). NAS owns the external -scheduler and its credentials. Wire contract: ``docs/chronos-managed-cron-contract.md``. -""" +"""Thin HTTP client for the agent -> NAS ``agent-cron`` endpoints (Chronos): arm one-shot / cancel / +list, authenticated with the existing Nous Portal token. +Wire contract: ``docs/chronos-managed-cron-contract.md``.""" from __future__ import annotations @@ -31,23 +28,20 @@ class NasCronClient: def _headers(self) -> Dict[str, str]: """Bearer auth with the agent's existing Nous Portal access token (refresh-aware).""" from hermes_cli.auth import resolve_nous_access_token - return { - "Authorization": f"Bearer {resolve_nous_access_token()}", - "Content-Type": "application/json"} + return {"Authorization": f"Bearer {resolve_nous_access_token()}", + "Content-Type": "application/json"} def _request(self, method: str, path: str, **kwargs: Any) -> Dict[str, Any]: """Issue one request; raise NasCronClientError on transport error or non-2xx.""" import requests # lazy: agent already depends on requests try: - resp = requests.request( - method, f"{self.portal_url}{path}", - headers=self._headers(), timeout=self.timeout_seconds, **kwargs) + resp = requests.request(method, f"{self.portal_url}{path}", headers=self._headers(), + timeout=self.timeout_seconds, **kwargs) except Exception as e: raise NasCronClientError(f"{method} {path} failed: {e}") from e if resp.status_code // 100 != 2: - raise NasCronClientError( - f"{method} {path} returned {resp.status_code}: {resp.text[:200]}") + raise NasCronClientError(f"{method} {path} returned {resp.status_code}: {resp.text[:200]}") try: return resp.json() if resp.content else {} except Exception: @@ -55,15 +49,10 @@ class NasCronClient: def provision(self, *, job_id: str, fire_at: str, agent_callback_url: str, dedup_key: str) -> Dict[str, Any]: - """Arm a one-shot for ``job_id`` at ``fire_at`` (ISO 8601). - - ``dedup_key`` (``{job_id}:{fire_at}``) makes re-arming the same fire idempotent - NAS-side. Returns the NAS response (e.g. ``{schedule_id}``). - """ + """Arm a one-shot for ``job_id`` at ``fire_at`` (ISO 8601); ``dedup_key`` makes re-arming + idempotent NAS-side. Returns the NAS response (e.g. ``{schedule_id}``).""" return self._request("POST", _PROVISION_PATH, json={ - "job_id": job_id, - "fire_at": fire_at, - "agent_callback_url": agent_callback_url, + "job_id": job_id, "fire_at": fire_at, "agent_callback_url": agent_callback_url, "dedup_key": dedup_key}) def cancel(self, *, job_id: str) -> Dict[str, Any]: diff --git a/plugins/cron_providers/chronos/verify.py b/plugins/cron_providers/chronos/verify.py index 67ae20ec6f..1b4feef1b2 100644 --- a/plugins/cron_providers/chronos/verify.py +++ b/plugins/cron_providers/chronos/verify.py @@ -1,10 +1,8 @@ """Inbound cron-fire token verification for Chronos. -NAS POSTs ``/api/cron/fire`` with a short-lived NAS-minted JWT; this verifies it (via -PyJWT, never hand-rolled) before any job runs. We verify a NAS-minted token rather than -let the external scheduler call the agent directly because the scheduler signs with NAS -keys the agent doesn't hold. ``get_fire_verifier`` is the pluggable seam so an -alternative auth mode (direct per-job cron-key) can swap in without a handler change. +NAS POSTs ``/api/cron/fire`` with a short-lived NAS-minted JWT, verified via PyJWT (never +hand-rolled) before any job runs. ``get_fire_verifier`` is the pluggable seam so an alternative +auth mode (direct per-job cron-key) can swap in without a handler change. """ from __future__ import annotations @@ -15,13 +13,11 @@ from typing import Any, Callable, Dict, Optional logger = logging.getLogger("cron.chronos.verify") -# Purpose claim scoping a token to the fire endpoint; a general agent JWT must not replay here. +# Purpose claim scoping a token to the fire endpoint (a general agent JWT must not replay here). _FIRE_PURPOSE = "cron_fire" -# Process-wide PyJWKClient cache keyed by JWKS URL. PyJWKClient caches keys on the -# INSTANCE, so a fresh client per fire re-fetched the JWKS every time — under bursts -# the portal rate-limited (403 -> 401) or the fetch blocked the loop past the relay -# timeout (504). One client per URL makes the steady state zero JWKS fetches per fire. +# Process-wide PyJWKClient cache keyed by JWKS URL: PyJWKClient caches keys on the INSTANCE, so a +# fresh client per fire re-fetched JWKS every time (portal rate-limits 403->401; relay 504s). _JWK_CLIENTS: Dict[str, Any] = {} _JWK_CLIENTS_LOCK = threading.Lock() @@ -43,24 +39,15 @@ def _get_jwk_client(jwks_url: str) -> Any: return client -def verify_nas_fire_token( - *, - token: str, - expected_audience: str, - jwks_or_key: Optional[str] = None, - issuer: Optional[str] = None, - leeway_seconds: int = 30) -> Optional[Dict[str, Any]]: - """Verify a NAS-minted cron-fire JWT; return decoded claims or None (never raises). - - Checks asymmetric signature (JWKS URL or inline PEM; symmetric secrets rejected since - NAS signs asymmetrically), ``aud``, ``exp``/``nbf`` with leeway, ``iss`` when configured, - and ``purpose == "cron_fire"``. Never raises, so the handler answers 401 without - leaking which check failed. - """ +def verify_nas_fire_token(*, token: str, expected_audience: str, jwks_or_key: Optional[str] = None, + issuer: Optional[str] = None, leeway_seconds: int = 30) -> Optional[Dict[str, Any]]: + """Verify a NAS-minted cron-fire JWT; return decoded claims or None (never raises, so the + handler answers 401 without leaking which check failed). Checks asymmetric signature (JWKS + URL or inline PEM; symmetric rejected), ``aud``, ``exp``/``nbf`` with leeway, ``iss`` when + configured, and ``purpose == "cron_fire"``.""" if not token or not expected_audience: return None - if not jwks_or_key: - # No key -> refuse; never fall back to unsigned decode on a security boundary. + if not jwks_or_key: # never fall back to unsigned decode on a security boundary logger.warning("cron fire: no JWKS/key configured; refusing token") return None @@ -73,10 +60,8 @@ def verify_nas_fire_token( signing_key = jwks_or_key # inline PEM public key (test / pinned-key deployments) decode_kwargs: Dict[str, Any] = dict( - algorithms=["RS256", "RS384", "RS512", "ES256", "ES384"], - audience=expected_audience, - leeway=leeway_seconds, - options={"require": ["exp", "aud"]}) + algorithms=["RS256", "RS384", "RS512", "ES256", "ES384"], audience=expected_audience, + leeway=leeway_seconds, options={"require": ["exp", "aud"]}) if issuer: decode_kwargs["issuer"] = issuer claims = jwt.decode(token, signing_key, **decode_kwargs) diff --git a/plugins/disk-cleanup/__init__.py b/plugins/disk-cleanup/__init__.py index 44938f3287..bdf7c8ba96 100644 --- a/plugins/disk-cleanup/__init__.py +++ b/plugins/disk-cleanup/__init__.py @@ -1,15 +1,8 @@ """disk-cleanup plugin — auto-cleanup of ephemeral Hermes session files. -Wires three behaviours: - -1. ``post_tool_call`` hook — inspects ``write_file`` / ``patch`` / ``terminal`` - tool calls for newly-created paths matching test/temp patterns under - ``HERMES_HOME`` and tracks them silently. Zero agent compliance required. -2. ``on_session_end`` hook — when any test files were auto-tracked during the - just-finished turn, runs :func:`disk_cleanup.quick` and logs one line to - ``$HERMES_HOME/disk-cleanup/cleanup.log``. -3. ``/disk-cleanup`` slash command — manual ``status``, ``dry-run``, ``quick``, - ``deep``, ``track``, ``forget``. +``post_tool_call`` silently tracks test/temp paths created by write_file/patch/terminal; +``on_session_end`` runs :func:`disk_cleanup.quick` when any test file was tracked this turn; +``/disk-cleanup`` exposes status / dry-run / quick / deep / track / forget. """ from __future__ import annotations @@ -26,27 +19,22 @@ from . import disk_cleanup as dg logger = logging.getLogger(__name__) -# Per-task set of test files newly tracked this turn, keyed by task_id (or -# session_id as fallback) so on_session_end can decide whether to run cleanup. -# Locked: post_tool_call can fire concurrently on parallel tool calls. +# Test files newly tracked this turn, keyed by task_id (or session_id) so on_session_end can +# decide whether to run cleanup. Locked: post_tool_call fires concurrently on parallel calls. _recent_test_tracks: Dict[str, Set[str]] = {} _lock = threading.Lock() _TERMINAL_PATH_REGEX = re.compile(r"(?:^|\s)(/[^\s'\"`]+|\~/[^\s'\"`]+)") -# --- Path extraction from tool calls ---------------------------------------- - def _extract_path_arg(args: Dict[str, Any], result: str) -> Set[str]: - """write_file and patch: only the single-file ``path`` arg. patch mostly edits - existing files; re-tracking is a no-op (track() dedups), so this is safe.""" + """write_file/patch: the single ``path`` arg (re-tracking existing files is a no-op).""" path = args.get("path") return {path} if isinstance(path, str) and path else set() def _extract_paths_from_terminal(args: Dict[str, Any], result: str) -> Set[str]: - """Best-effort: pull candidate filesystem paths from a terminal command and - its output; ``guess_category`` / ``is_safe_path`` filter them afterwards.""" + """Candidate paths from a terminal command + output; guess_category/is_safe_path filter later.""" paths: Set[str] = set() cmd = args.get("command") or "" if isinstance(cmd, str) and cmd: @@ -66,16 +54,8 @@ _PATH_EXTRACTORS: Dict[str, Callable[[Dict[str, Any], str], Set[str]]] = { "terminal": _extract_paths_from_terminal} -# --- Hooks ------------------------------------------------------------------ - -def _on_post_tool_call( - tool_name: str = "", - args: Optional[Dict[str, Any]] = None, - result: Any = None, - task_id: str = "", - session_id: str = "", - tool_call_id: str = "", - **_: Any) -> None: +def _on_post_tool_call(tool_name: str = "", args: Optional[Dict[str, Any]] = None, result: Any = None, + task_id: str = "", session_id: str = "", tool_call_id: str = "", **_: Any) -> None: """Auto-track ephemeral files created by recent tool calls. Best-effort, never raises.""" extractor = _PATH_EXTRACTORS.get(tool_name) if not isinstance(args, dict) or extractor is None: @@ -96,8 +76,7 @@ def _on_post_tool_call( def _on_session_end( session_id: str = "", completed: bool = True, interrupted: bool = False, **_: Any) -> None: """Run quick cleanup if any test files were tracked during this turn.""" - # Drain the session bucket plus every task-scoped bucket: subagents record - # into their own task_id buckets and should be cleaned up at session end too. + # Drain the session bucket plus every task-scoped bucket (subagents record into their own). with _lock: had_tracks = bool(_recent_test_tracks.pop(session_id or "default", None) or _recent_test_tracks) _recent_test_tracks.clear() @@ -111,13 +90,10 @@ def _on_session_end( return if summary["deleted"] or summary["empty_dirs"]: - dg._log( - f"AUTO_QUICK (session_end): deleted={summary['deleted']} " - f"dirs={summary['empty_dirs']} freed={dg.fmt_size(summary['freed'])}") + dg._log(f"AUTO_QUICK (session_end): deleted={summary['deleted']} " + f"dirs={summary['empty_dirs']} freed={dg.fmt_size(summary['freed'])}") -# --- Slash command ---------------------------------------------------------- - _HELP_TEXT = """\ /disk-cleanup — ephemeral-file cleanup @@ -137,9 +113,8 @@ Test files are auto-tracked on write_file / terminal and auto-cleaned at session def _fmt_summary(summary: Dict[str, Any]) -> str: - base = ( - f"[disk-cleanup] Cleaned {summary['deleted']} files + " - f"{summary['empty_dirs']} empty dirs, freed {dg.fmt_size(summary['freed'])}.") + base = (f"[disk-cleanup] Cleaned {summary['deleted']} files + " + f"{summary['empty_dirs']} empty dirs, freed {dg.fmt_size(summary['freed'])}.") if summary.get("errors"): base += f"\n {len(summary['errors'])} error(s); see cleanup.log." return base @@ -161,8 +136,7 @@ def _cmd_dry_run(argv: List[str]) -> str: def _cmd_deep(argv: List[str]) -> str: - # In-session deep can't prompt interactively — show what quick cleaned plus - # the items that WOULD need confirmation. + # In-session deep can't prompt — show what quick cleaned plus items needing confirmation. quick_summary = dg.quick() _auto, prompt_items = dg.dry_run() lines = [_fmt_summary(quick_summary)] @@ -187,9 +161,8 @@ def _cmd_forget(argv: List[str]) -> str: if len(argv) < 2: return "Usage: /disk-cleanup forget " n = dg.forget(argv[1]) - return ( - f"Removed {n} tracking entr{'y' if n == 1 else 'ies'} for {argv[1]}." - if n else f"Not found in tracking: {argv[1]}") + return (f"Removed {n} tracking entr{'y' if n == 1 else 'ies'} for {argv[1]}." if n + else f"Not found in tracking: {argv[1]}") _SUBCOMMANDS: Dict[str, Callable[[List[str]], str]] = { @@ -214,7 +187,5 @@ def _handle_slash(raw_args: str) -> Optional[str]: def register(ctx) -> None: ctx.register_hook("post_tool_call", _on_post_tool_call) ctx.register_hook("on_session_end", _on_session_end) - ctx.register_command( - "disk-cleanup", - handler=_handle_slash, - description="Track and clean up ephemeral Hermes session files.") + ctx.register_command("disk-cleanup", handler=_handle_slash, + description="Track and clean up ephemeral Hermes session files.") diff --git a/plugins/disk-cleanup/disk_cleanup.py b/plugins/disk-cleanup/disk_cleanup.py index 884c5aecb2..0ec670fe35 100755 --- a/plugins/disk-cleanup/disk_cleanup.py +++ b/plugins/disk-cleanup/disk_cleanup.py @@ -1,15 +1,9 @@ -"""disk_cleanup — ephemeral file cleanup for Hermes Agent. - -Library module behind the disk-cleanup plugin; ``__init__.py`` wires these -functions into ``post_tool_call`` / ``on_session_end`` hooks so tracking and -cleanup happen without the agent calling a tool or remembering a skill. +"""disk_cleanup — ephemeral file cleanup library behind the disk-cleanup plugin. Rules: test files delete at task end (age >= 0); temp after 7 days; cron-output -after 14 days; empty dirs under HERMES_HOME always. Deep-only prompts: research +after 14 days; empty dirs under HERMES_HOME always. Prompt-only: research (keep 10 newest, > 30 days), chrome-profile > 14 days, any file > 500 MB. - -Scope: strictly HERMES_HOME and /tmp/hermes-*. Never touches ~/.hermes/logs/ or -any system directory. +Scope: strictly HERMES_HOME and /tmp/hermes-*; never ~/.hermes/logs/ or system dirs. """ from __future__ import annotations @@ -19,7 +13,7 @@ import logging import shutil from datetime import datetime, timezone from pathlib import Path -from typing import Any, Callable, Dict, Iterator, List, Optional, Tuple +from typing import Any, Dict, Iterator, List, Optional, Tuple try: from hermes_constants import get_hermes_home @@ -36,19 +30,13 @@ logger = logging.getLogger(__name__) _LARGE_FILE_BYTES = 500 * 1024 * 1024 -# --- Paths / safety --------------------------------------------------------- - def _state_file(name: str) -> Path: - """``$HERMES_HOME/disk-cleanup/`` — state and audit log deliberately - live outside ``$HERMES_HOME/logs/``.""" + """``$HERMES_HOME/disk-cleanup/`` — deliberately outside ``$HERMES_HOME/logs/``.""" return get_hermes_home() / "disk-cleanup" / name def is_safe_path(path: Path) -> bool: - """Accept only paths under HERMES_HOME or ``/tmp/hermes-*``. - - Rejects Windows mounts (``/mnt/c`` etc.) and any system directory. - """ + """Accept only paths under HERMES_HOME or ``/tmp/hermes-*`` (rejects /mnt/c etc.).""" try: path.resolve().relative_to(get_hermes_home()) return True @@ -70,8 +58,6 @@ def _log(message: str) -> None: pass -# --- tracked.json — atomic read/write, backup scoped to tracked.json only ---- - def load_tracked() -> List[Dict[str, Any]]: """Load tracked.json. Restores from ``.bak`` on corruption.""" tf = _state_file("tracked.json") @@ -81,16 +67,17 @@ def load_tracked() -> List[Dict[str, Any]]: try: return json.loads(tf.read_text(encoding="utf-8")) except (json.JSONDecodeError, ValueError): - bak = tf.with_suffix(".json.bak") - if bak.exists(): - try: - data = json.loads(bak.read_text(encoding="utf-8")) - _log("WARN: tracked.json corrupted — restored from .bak") - return data - except Exception: - pass - _log("WARN: tracked.json corrupted, no backup — starting fresh") - return [] + pass + bak = tf.with_suffix(".json.bak") + if bak.exists(): + try: + data = json.loads(bak.read_text(encoding="utf-8")) + _log("WARN: tracked.json corrupted — restored from .bak") + return data + except Exception: + pass + _log("WARN: tracked.json corrupted, no backup — starting fresh") + return [] def save_tracked(tracked: List[Dict[str, Any]]) -> None: @@ -104,13 +91,10 @@ def save_tracked(tracked: List[Dict[str, Any]]) -> None: tmp.replace(tf) -# --- Categories / protected trees ------------------------------------------- - ALLOWED_CATEGORIES = { "temp", "test", "research", "download", "chrome-profile", "cron-output", "other"} -# Top-level HERMES_HOME dirs whose empty subdirs are never swept. The last row -# is user-authored project trees (patches/, projects/, ...) — never sweep inside. +# Top-level HERMES_HOME dirs whose empty subdirs are never swept (last row: user project trees). _EMPTY_DIR_PROTECTED_TOP_LEVEL = frozenset({ "logs", "memories", "sessions", "cron", "cronjobs", "cache", "skills", "plugins", "disk-cleanup", "optional-skills", @@ -120,9 +104,8 @@ _EMPTY_DIR_PROTECTED_TOP_LEVEL = frozenset({ _EMPTY_DIR_SWEEP_PRUNE_DIRS = frozenset({ ".git", "node_modules", "venv", ".venv", "site-packages", "__pycache__"}) -# Top-level entries under HERMES_HOME that guess_category() never auto-tracks: -# state dir, logs, memory, sessions, config/secrets, and user-authored project -# trees (a file named test_*/tmp_* inside patches/ or projects/ is not disposable). +# Top-level HERMES_HOME entries guess_category() never auto-tracks: state, logs, memory, +# sessions, config/secrets, and user project trees (test_* inside projects/ is not disposable). _NEVER_TRACK_TOP_LEVEL = frozenset({ "disk-cleanup", "logs", "memories", "sessions", "config.yaml", "skills", "plugins", ".env", "USER.md", "MEMORY.md", "SOUL.md", @@ -130,20 +113,15 @@ _NEVER_TRACK_TOP_LEVEL = frozenset({ "patches", "projects", "skins", "themes", "contributors", "profiles", "backups", "optional-skills"}) -# Defense-in-depth for quick(): exact cron control-plane paths never deleted, -# regardless of stored category (guards stale tracked.json entries). +# Defense-in-depth for quick(): exact cron control-plane paths never deleted regardless of +# stored category (guards stale tracked.json entries). _PROTECTED_CRON_PATHS: set[str] = set() def _is_protected_cron_path(p: Path) -> bool: - """True if *p* is cron control-plane state that must never be deleted. - - Matches by EXACT path only: the ``cron/`` dir itself, ``jobs.json``, - ``.tick.lock``, and the ``output/`` root. It must NOT be widened to - everything under ``cron/output/`` — run artifacts there are disposable and - cleaned by retention; only the ``output/`` root is protected because - deleting it wholesale erases every job's retained run history. - """ + """True if *p* is cron control-plane state (EXACT match: ``cron/``, ``jobs.json``, + ``.tick.lock``, the ``output/`` root). Never widen to everything under ``cron/output/``: + run artifacts there are disposable; only wholesale deletion of ``output/`` is fatal.""" if not _PROTECTED_CRON_PATHS: # built lazily so HERMES_HOME resolves once hermes_home = get_hermes_home() for parent in ("cron", "cronjobs"): @@ -161,8 +139,6 @@ def fmt_size(n: float) -> str: return f"{n:.1f} PB" -# --- Track / forget --------------------------------------------------------- - def track(path_str: str, category: str, silent: bool = False) -> bool: """Register a file for tracking. Returns True if newly tracked.""" if category not in ALLOWED_CATEGORIES: @@ -182,11 +158,8 @@ def track(path_str: str, category: str, silent: bool = False) -> bool: if any(item["path"] == str(path) for item in tracked): return False - tracked.append({ - "path": str(path), - "timestamp": datetime.now(timezone.utc).isoformat(), - "category": category, - "size": size}) + tracked.append({"path": str(path), "timestamp": datetime.now(timezone.utc).isoformat(), + "category": category, "size": size}) save_tracked(tracked) _log(f"TRACKED: {path} ({category}, {fmt_size(size)})") if not silent: @@ -198,26 +171,22 @@ def forget(path_str: str) -> int: """Remove a path from tracking without deleting the file.""" p = Path(path_str).resolve() tracked = load_tracked() - before = len(tracked) - tracked = [i for i in tracked if Path(i["path"]).resolve() != p] - removed = before - len(tracked) + kept = [i for i in tracked if Path(i["path"]).resolve() != p] + removed = len(tracked) - len(kept) if removed: - save_tracked(tracked) + save_tracked(kept) _log(f"FORGOT: {p} ({removed} entries)") return removed -# --- Rules shared by dry_run / quick / deep --------------------------------- - def _live_items(tracked: List[Dict], now: datetime, *, log_stale: bool = False) -> Iterator[Tuple[Dict, Path, int]]: """Yield ``(item, path, age_days)`` for entries whose path still exists.""" for item in tracked: p = Path(item["path"]) - if not p.exists(): - if log_stale: - _log(f"STALE: {p} (removed from tracking)") - continue - yield item, p, (now - datetime.fromisoformat(item["timestamp"])).days + if p.exists(): + yield item, p, (now - datetime.fromisoformat(item["timestamp"])).days + elif log_stale: + _log(f"STALE: {p} (removed from tracking)") def _is_auto_delete(cat: str, age: int) -> bool: @@ -225,15 +194,13 @@ def _is_auto_delete(cat: str, age: int) -> bool: def _prompt_group(item: Dict, age: int) -> Optional[str]: - """Deep-only bucket: ``research`` / ``chrome`` / ``large`` or None.""" + """Prompt-only bucket: ``research`` / ``chrome`` / ``large`` or None.""" cat = item["category"] if cat == "research" and age > 30: return "research" if cat == "chrome-profile" and age > 14: return "chrome" - if item["size"] > _LARGE_FILE_BYTES: - return "large" - return None + return "large" if item["size"] > _LARGE_FILE_BYTES else None def _delete_item(item: Dict) -> Optional[str]: @@ -251,15 +218,11 @@ def _delete_item(item: Dict) -> Optional[str]: return None -# Stored categories that are re-validated against guess_category() before use. -# Old tracked.json entries can carry "cron-output" for control-plane files -# (cron/jobs.json) or "test" for files now under protected project trees; -# guess_category() was tightened later but existing entries were never re-checked. +# Stored categories re-validated against guess_category() before use: old tracked.json entries +# may carry "cron-output" for control-plane files or "test" for files under protected trees. _STALE_SKIP_NOTE = {"cron-output": "", "test": " — under protected tree"} -# --- Dry run / quick / deep ------------------------------------------------- - def dry_run() -> Tuple[List[Dict], List[Dict]]: """Return (auto_delete_list, needs_prompt_list) without touching files.""" auto: List[Dict] = [] @@ -277,10 +240,7 @@ def dry_run() -> Tuple[List[Dict], List[Dict]]: def quick() -> Dict[str, Any]: - """Safe deterministic cleanup — no prompts. - - Returns: ``{"deleted": N, "empty_dirs": N, "freed": bytes, "errors": [str, ...]}``. - """ + """Safe deterministic cleanup — no prompts. Returns ``{deleted, empty_dirs, freed, errors}``.""" deleted = freed = 0 new_tracked: List[Dict] = [] errors: List[str] = [] @@ -320,9 +280,8 @@ def _subdirs(dirpath: Path, exclude: frozenset) -> List[Path]: def _sweep_empty_dirs(hermes_home: Path) -> int: - """Remove empty dirs under HERMES_HOME, never recursing into durable state - trees. Some installs keep the Hermes checkout, venv, and desktop build under - HERMES_HOME; a full rglob there can stall the gateway event loop for minutes. + """Remove empty dirs under HERMES_HOME without recursing into durable/heavy trees (a full + rglob over a checkout+venv under HERMES_HOME can stall the gateway loop for minutes). Iterative post-order so parents emptied by child removal are caught.""" removed = 0 stack: List[Tuple[Path, bool]] = [ @@ -344,36 +303,6 @@ def _sweep_empty_dirs(hermes_home: Path) -> int: return removed -def deep(confirm: Optional[Callable[[Dict], bool]] = None) -> Dict[str, Any]: - """Deep cleanup: :func:`quick`, then ask *confirm(item)* for each risky item - (research > 30d beyond the 10 newest, chrome-profile > 14d, any file > 500 MB). - - Returns: ``{"quick": {...}, "deep_deleted": N, "deep_freed": bytes}``. - """ - quick_result = quick() - if confirm is None: # no interactive confirmer — stop after the quick pass - return {"quick": quick_result, "deep_deleted": 0, "deep_freed": 0} - - tracked = load_tracked() - groups: Dict[str, List[Dict]] = {"research": [], "chrome": [], "large": []} - for item, _p, age in _live_items(tracked, datetime.now(timezone.utc)): - group = _prompt_group(item, age) - if group: - groups[group].append(item) - - groups["research"].sort(key=lambda x: x["timestamp"], reverse=True) - del groups["research"][:10] # keep the 10 newest research items - - removed = [item for group in groups.values() for item in group if confirm(item) and _delete_item(item) is None] - if removed: - remove_paths = {i["path"] for i in removed} - save_tracked([i for i in tracked if i["path"] not in remove_paths]) - - return {"quick": quick_result, "deep_deleted": len(removed), "deep_freed": sum(i["size"] for i in removed)} - - -# --- Status ----------------------------------------------------------------- - def status() -> Dict[str, Any]: """Return per-category breakdown and top 10 largest tracked files.""" tracked = load_tracked() @@ -383,8 +312,8 @@ def status() -> Dict[str, Any]: c["count"] += 1 c["size"] += item["size"] - existing = [(i["path"], i["size"], i["category"]) for i in tracked if Path(i["path"]).exists()] - existing.sort(key=lambda x: x[1], reverse=True) + existing = sorted(((i["path"], i["size"], i["category"]) for i in tracked + if Path(i["path"]).exists()), key=lambda x: x[1], reverse=True) return {"categories": cats, "top10": existing[:10], "total_tracked": len(tracked)} @@ -405,17 +334,12 @@ def format_status(s: Dict[str, Any]) -> str: return "\n".join(lines) -# --- Auto-categorisation from tool-call inspection -------------------------- - _TEST_PATTERNS = ("test_", "tmp_") _TEST_SUFFIXES = (".test.py", ".test.js", ".test.ts", ".test.md") def guess_category(path: Path) -> Optional[str]: - """Return a category label for *path*, or None if we shouldn't track it. - - Used by the ``post_tool_call`` hook to auto-track ephemeral files. - """ + """Category label for *path*, or None if we shouldn't track it (``post_tool_call`` hook).""" if not is_safe_path(path): return None @@ -425,18 +349,13 @@ def guess_category(path: Path) -> Optional[str]: if top in _NEVER_TRACK_TOP_LEVEL: return None if top in ("cron", "cronjobs"): - # Only the disposable ``output/`` subtree is a candidate. Top-level - # control-plane state (jobs.json, .tick.lock) must never be tracked — - # deleting it wipes the live scheduler registry. - if len(rel.parts) >= 3 and rel.parts[1] == "output": - return "cron-output" - return None + # Only the disposable ``output/`` subtree; control-plane state (jobs.json, + # .tick.lock) must never be tracked — deleting it wipes the scheduler registry. + return "cron-output" if len(rel.parts) >= 3 and rel.parts[1] == "output" else None if top == "cache": return "temp" except ValueError: pass # not under HERMES_HOME (e.g. /tmp/hermes-*) — fall through to name rules name = path.name - if name.startswith(_TEST_PATTERNS) or name.endswith(_TEST_SUFFIXES): - return "test" - return None + return "test" if name.startswith(_TEST_PATTERNS) or name.endswith(_TEST_SUFFIXES) else None diff --git a/plugins/google_meet/__init__.py b/plugins/google_meet/__init__.py index faacddedf1..19e12cd55b 100644 --- a/plugins/google_meet/__init__.py +++ b/plugins/google_meet/__init__.py @@ -1,12 +1,9 @@ """google_meet plugin — let the agent join a Meet call, transcribe it, follow up. -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. +Headless Chromium (Playwright) joins the URL, enables live captions and scrapes them into a +transcript. Realtime mode adds agent speech (OpenAI Realtime + virtual audio device); remote +nodes run the bot on another machine. Explicit-by-design: only joins ``https://meet.google.com/`` +URLs passed in — no calendar scanning, auto-dial or consent announcement. """ from __future__ import annotations @@ -34,10 +31,7 @@ _TOOLS = ( def _on_session_end(**kwargs) -> None: - """Leave a still-running call so we don't orphan a headless Chromium. - - Swallows all exceptions — session end must never fail on bot cleanup. - """ + """Leave a still-running call so we don't orphan a headless Chromium (never raises).""" try: status = pm.status() if status.get("ok") and status.get("alive"): @@ -48,8 +42,7 @@ def _on_session_end(**kwargs) -> None: def register(ctx) -> None: """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. + # Windows: no tested audio-routing path and flaky guest-join Chromium — refuse 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) @@ -60,12 +53,9 @@ def register(ctx) -> None: check_fn=check_meet_requirements, emoji=emoji) ctx.register_cli_command( - name="meet", - help="Google Meet bot (join, transcribe, follow up)", - setup_fn=_register_meet_cli, - handler_fn=_meet_command, - description=( - "Let the hermes agent join a Google Meet call and scrape live " - "captions into a transcript. See: hermes meet setup")) + name="meet", help="Google Meet bot (join, transcribe, follow up)", + setup_fn=_register_meet_cli, handler_fn=_meet_command, + description=("Let the hermes agent join a Google Meet call and scrape live " + "captions into a transcript. See: hermes meet setup")) ctx.register_hook("on_session_end", _on_session_end) diff --git a/plugins/google_meet/_jsonfile.py b/plugins/google_meet/_jsonfile.py index 24a36a1a9a..9758e57e7e 100644 --- a/plugins/google_meet/_jsonfile.py +++ b/plugins/google_meet/_jsonfile.py @@ -18,11 +18,8 @@ def read_json(path: Path) -> Optional[Any]: 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. - """ + """Write ``json.dumps(data, indent=2)`` via a ``.json.tmp`` sibling + rename; *mode* (e.g. + ``0o600``) is applied to the temp file so the final file never exists with looser perms.""" 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") diff --git a/plugins/google_meet/audio_bridge.py b/plugins/google_meet/audio_bridge.py index ad3b7046db..1ad6a2ce03 100644 --- a/plugins/google_meet/audio_bridge.py +++ b/plugins/google_meet/audio_bridge.py @@ -1,14 +1,8 @@ """Virtual audio bridge for feeding generated speech into Chrome's mic. -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: 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: 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. +Linux: pactl creates a null-sink plus a virtual source on the sink's monitor; callers set +``PULSE_SOURCE=`` in Chrome's env. macOS: only verifies BlackHole 2ch is installed +(the default-input switch is left to the user). Windows: unsupported. """ from __future__ import annotations @@ -22,20 +16,12 @@ _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) + 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 (idempotent). - """ + """Virtual audio device for Chrome fake-mic input: ``setup()`` before launch, ``teardown()`` after.""" def __init__(self, name_prefix: str = "hermes_meet") -> None: self._name_prefix = name_prefix @@ -69,8 +55,7 @@ class AudioBridge: if self._torn_down: return if self._platform == "linux": - # Unload in reverse order (virtual-source before null-sink). - for mod_id in reversed(self._module_ids): + for mod_id in reversed(self._module_ids): # virtual-source before null-sink try: _pactl("unload-module", str(mod_id), check=False) except Exception: @@ -114,19 +99,16 @@ class AudioBridge: def _setup_darwin(self) -> dict: try: - out = subprocess.check_output( - ["system_profiler", "SPAudioDataType"], - text=True, encoding='utf-8', errors='replace', - stderr=subprocess.STDOUT) + out = subprocess.check_output(["system_profiler", "SPAudioDataType"], text=True, + encoding='utf-8', errors='replace', stderr=subprocess.STDOUT) except FileNotFoundError as 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 if "BlackHole" not in out: - raise RuntimeError( - "BlackHole virtual audio device not installed. " - "Install via: brew install blackhole-2ch") + raise RuntimeError("BlackHole virtual audio device not installed. " + "Install via: brew install blackhole-2ch") return self._finish("darwin", _BLACKHOLE_DEVICE, _BLACKHOLE_DEVICE, []) @staticmethod diff --git a/plugins/google_meet/cli.py b/plugins/google_meet/cli.py index c8b7a22e94..f10ee898de 100644 --- a/plugins/google_meet/cli.py +++ b/plugins/google_meet/cli.py @@ -22,6 +22,7 @@ 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.node.cli import register_cli as _register_node_cli from plugins.google_meet.tools import resolve_node @@ -37,14 +38,10 @@ def register_cli(subparser: argparse.ArgumentParser) -> None: inst_p = subs.add_parser( "install", help="Install prerequisites (pip deps, Chromium, platform audio tools)") - inst_p.add_argument( - "--realtime", action="store_true", - help="Also install realtime audio tools (pulseaudio-utils on Linux, BlackHole+ffmpeg on macOS). Uses sudo/brew, prompts before invoking either.", - ) - inst_p.add_argument( - "--yes", "-y", action="store_true", - help="Answer yes to all prompts (use with care; will run sudo apt-get or brew without asking).", - ) + inst_p.add_argument("--realtime", action="store_true", + help="Also install realtime audio tools (pulseaudio-utils on Linux, BlackHole+ffmpeg on macOS). Uses sudo/brew, prompts before invoking either.") + inst_p.add_argument("--yes", "-y", action="store_true", + help="Answer yes to all prompts (use with care; will run sudo apt-get or brew without asking).") subs.add_parser("auth", help="Sign in to Google and save session state") @@ -53,11 +50,10 @@ def register_cli(subparser: argparse.ArgumentParser) -> None: join_p.add_argument("--guest-name", default="Hermes Agent") join_p.add_argument("--duration", default=None, help="e.g. 30m, 2h, 90s") join_p.add_argument("--headed", action="store_true", help="show browser") - join_p.add_argument( - "--mode", choices=("transcribe", "realtime"), default="transcribe", - help="transcribe (default, listen-only) or realtime (speak via OpenAI Realtime)") - join_p.add_argument( - "--node", default=None, help="remote node name, or 'auto' to use the sole registered node") + join_p.add_argument("--mode", choices=("transcribe", "realtime"), default="transcribe", + help="transcribe (default, listen-only) or realtime (speak via OpenAI Realtime)") + join_p.add_argument("--node", default=None, + help="remote node name, or 'auto' to use the sole registered node") subs.add_parser("status", help="Print current Meet bot state") @@ -70,20 +66,8 @@ def register_cli(subparser: argparse.ArgumentParser) -> None: subs.add_parser("stop", help="Leave the current meeting") - node_p = subs.add_parser( - "node", help="Manage remote meet node hosts (run/list/approve/remove/status/ping)") - try: - 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 - # Keep the subparser present so argparse dispatch surfaces a clear error. - err = e - - def _node_unavailable(args): - print(f"hermes meet node: module unavailable ({err})") - return 1 - node_p.set_defaults(func=_node_unavailable) - + _register_node_cli(subs.add_parser( + "node", help="Manage remote meet node hosts (run/list/approve/remove/status/ping)")) subparser.set_defaults(func=meet_command) @@ -97,12 +81,11 @@ def _cmd_node(args: argparse.Namespace) -> int: _DISPATCH = { "setup": lambda a: _cmd_setup(), - "install": lambda a: _cmd_install( - realtime=bool(getattr(a, "realtime", False)), assume_yes=bool(getattr(a, "yes", False))), + "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)), + "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)), @@ -155,24 +138,19 @@ def _cmd_setup() -> int: print(f" chromium : {chromium_msg}") auth_path = _auth_state_path() - print( - " google auth : " - + (f"ok ({auth_path})" if auth_path.is_file() else "not saved — run: hermes meet auth")) + print(" google 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") - else: - print("not ready yet — fix the items above.") + print("ready. Join a meeting: hermes meet join https://meet.google.com/abc-defg-hij" if all_ok + else "not ready yet — fix the items above.") return 0 if all_ok else 1 def _cmd_install(*, realtime: bool, assume_yes: bool) -> int: - """pip deps + Chromium; with ``--realtime`` also the platform audio bridge deps. - - Prompts before every package-manager invocation unless ``--yes``. Linux/macOS only. - """ + """pip deps + Chromium; ``--realtime`` adds the platform audio bridge deps. + Prompts before every package-manager invocation unless ``--yes``. Linux/macOS only.""" system = platform.system() if system not in {"Linux", "Darwin"}: print(f"google_meet install: {system} is not supported (linux/macos only)") @@ -242,17 +220,15 @@ def _cmd_install(*, realtime: bool, assume_yes: bool) -> int: if not needs: print(" BlackHole and ffmpeg already installed.") elif not shutil.which("brew"): - print( - " missing: " + ", ".join(needs) + "\n" - " install Homebrew first (https://brew.sh) or install the packages manually.") + print(" missing: " + ", ".join(needs) + "\n" + " install Homebrew first (https://brew.sh) or install the packages manually.") else: _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.") + 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.") print("\ndone. verify with: hermes meet setup") return 0 @@ -263,9 +239,8 @@ def _cmd_auth() -> int: try: from playwright.sync_api import sync_playwright except ImportError: - print( - "playwright is not installed. run:\n" - " pip install playwright && python -m playwright install chromium") + print("playwright is not installed. run:\n" + " pip install playwright && python -m playwright install chromium") return 1 path = _auth_state_path() @@ -315,14 +290,8 @@ def _remote(node: str, op: str, call) -> int: return _print_result({"node": name, **res}) -def _cmd_join( - url: str, - *, - guest_name: str, - duration: Optional[str], - headed: bool, - mode: str = "transcribe", - node: Optional[str] = None) -> int: +def _cmd_join(url: str, *, guest_name: str, duration: Optional[str], headed: bool, + mode: str = "transcribe", node: Optional[str] = None) -> int: if not _is_safe_meet_url(url): print(f"refusing: not a meet.google.com URL: {url}") return 2 @@ -330,9 +299,8 @@ def _cmd_join( 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() - 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)) + 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: diff --git a/plugins/google_meet/meet_bot.py b/plugins/google_meet/meet_bot.py index a9852e2941..c84ef0e8df 100644 --- a/plugins/google_meet/meet_bot.py +++ b/plugins/google_meet/meet_bot.py @@ -1,16 +1,11 @@ """Headless Google Meet bot — Playwright + live-caption scraping. -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. - -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. - -Debug run: ``HERMES_MEET_URL=... HERMES_MEET_OUT_DIR=/tmp/meet-debug -HERMES_MEET_HEADED=1 python -m plugins.google_meet.meet_bot`` +Standalone subprocess spawned by ``process_manager.py``; configured via ``HERMES_MEET_*`` env, +status + transcript written under ``$HERMES_MEET_OUT_DIR`` (filesystem is the only IPC). +No WebRTC audio parsing: Meet's live captions are watched via a MutationObserver — lossy and +English-biased, but deterministic (no STT billing) and stable thanks to the ARIA role. +Debug: ``HERMES_MEET_URL=... HERMES_MEET_OUT_DIR=/tmp/x HERMES_MEET_HEADED=1 \\ + python -m plugins.google_meet.meet_bot`` """ from __future__ import annotations @@ -31,15 +26,8 @@ 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,}" - r"|lookup/[^/?#]+" - r"|new" - r")(?:[/?#].*)?$") - -# Filenames the bot reads/writes in ``HERMES_MEET_OUT_DIR``. -SAY_QUEUE_FILENAME = "say_queue.jsonl" -SAY_PCM_FILENAME = "speaker.pcm" + r"^https://meet\.google\.com/([a-z0-9]{3,}-[a-z0-9]{3,}-[a-z0-9]{3,}|lookup/[^/?#]+|new)" + r"(?:[/?#].*)?$") _FFMPEG_MISSING = "ffmpeg not found — install via `brew install ffmpeg` for realtime on macOS" @@ -65,15 +53,13 @@ def _quiet(fn, *args, **kwargs): # 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), + ("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 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), @@ -113,8 +99,7 @@ class _BotState: def _flush(self) -> None: 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 + data.update(transcriptPath=str(self.transcript_path), pid=os.getpid()) # keeps table key order write_json_atomic(self.status_path, data) def set(self, **kwargs) -> None: @@ -182,8 +167,7 @@ _CAPTION_OBSERVER_JS = r""" })(); """ -# Best-effort caption toggle: Meet binds it to the ``c`` key; click targeting -# is too brittle to rely on. +# Best-effort caption toggle: Meet binds it to the ``c`` key; click targeting is too brittle. _ENABLE_CAPTIONS_JS = ( "(() => { document.body.dispatchEvent(new KeyboardEvent('keydown', " "{ key: 'c', code: 'KeyC', keyCode: 67, which: 67, bubbles: true })); return true; })();") @@ -192,8 +176,8 @@ _LEAVE_CALL_JS = ( "() => { const b = document.querySelector('button[aria-label*=\"eave call\"]');" " if (b) b.click(); }") -# True once we're clearly past the lobby: leave button, caption region -# (only once our observer is installed) or participant list visible. +# True once past the lobby: leave button, caption region (once our observer is installed) +# or participant list visible. _ADMISSION_PROBE_JS = r""" (() => { if (document.querySelector('button[aria-label*="eave call" i]')) return true; @@ -230,7 +214,7 @@ def _visible(locator): 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.""" + """Stream the growing ``speaker.pcm`` (24kHz s16le mono) into the device Chrome's fake mic reads.""" bridge_info = bridge_info or {} platform_tag = bridge_info.get("platform") if platform_tag == "linux": @@ -239,8 +223,7 @@ def _start_pcm_pump(rt: dict, bridge_info: dict, pcm_path: Path, state: "_BotSta 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. + # User must have BlackHole as default input; ffmpeg targets it by audiotoolbox index. if not shutil.which("ffmpeg"): state.set(error=_FFMPEG_MISSING) return @@ -270,9 +253,9 @@ def _start_realtime_speaker(rt: dict, cfg: "_BotConfig", stop_flag: dict, state: state.set(error=f"realtime import failed: {e}") return - 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 + pcm_path = cfg.out_dir / "speaker.pcm" + queue_path = cfg.out_dir / "say_queue.jsonl" + pcm_path.write_bytes(b"") # clean sink file per session queue_path.touch() # so the speaker poller doesn't error on first iteration try: @@ -301,10 +284,8 @@ def _start_realtime_speaker(rt: dict, cfg: "_BotConfig", stop_flag: dict, state: def _mac_audio_device_index(device_name: str) -> str: - """ffmpeg ``-audio_device_index`` for *device_name* (case-insensitive), ``"0"`` if not found. - - ffmpeg prints the avfoundation device table on stderr as ``[N] Name``. - """ + """ffmpeg ``-audio_device_index`` for *device_name* (case-insensitive; ``"0"`` if not found). + ffmpeg prints the avfoundation device table on stderr as ``[N] Name``.""" try: out = subprocess.run( ["ffmpeg", "-f", "avfoundation", "-list_devices", "true", "-i", ""], @@ -342,10 +323,9 @@ def _teardown_realtime(rt: dict) -> None: _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) + for key, method in (("session", "close"), ("bridge", "teardown")): + if rt[key]: + _quiet(getattr(rt[key], method)) @dataclass @@ -377,9 +357,8 @@ def _config_from_env() -> _BotConfig: 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. + # HERMES_MEET_REALTIME_KEY is resolved by process_manager.start() via the parent's + # profile secret scope; OPENAI_API_KEY only serves standalone `python -m` 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"), @@ -388,10 +367,7 @@ def _config_from_env() -> _BotConfig: 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``. - """ + """Fill the guest-name field and click 'Join now' / 'Ask to join' (the latter → 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) @@ -409,11 +385,8 @@ def _join(page, cfg: _BotConfig, state: _BotState) -> None: 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. - """ + """Admission + caption drain loop until SIGTERM, duration expiry, lobby timeout/denial or page loss. + Sets ``leave_reason`` for every exit but SIGTERM; triggers barge-in; mirrors realtime counters.""" deadline = (time.time() + cfg.duration_s) if cfg.duration_s else None lobby_deadline = time.time() + cfg.lobby_timeout last_admission_check = 0.0 @@ -429,9 +402,8 @@ def _drain_loop(page, cfg: _BotConfig, state: _BotState, rt: dict, stop_flag: di 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") + 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") @@ -444,22 +416,18 @@ def _drain_loop(page, cfg: _BotConfig, state: _BotState, rt: dict, stop_flag: di 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) + # Barge-in: a real human spoke while we may be generating audio. + 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(): + if page.is_closed(): # Meet reloaded or we got booted — exit rather than spin 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) + state.set(audio_bytes_out=rt["session"].audio_bytes_out, + last_audio_out_at=rt["session"].last_audio_out_at) time.sleep(1.0) @@ -467,9 +435,8 @@ def _drain_loop(page, cfg: _BotConfig, state: _BotState, rt: dict, stop_flag: di 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" % cfg.url) + sys.stderr.write("google_meet bot: refusing to launch — HERMES_MEET_URL must be a " + "meet.google.com URL. got: %r\n" % cfg.url) return 2 if cfg.out_dir is None: sys.stderr.write("google_meet bot: HERMES_MEET_OUT_DIR is required\n") @@ -477,19 +444,15 @@ def run_bot() -> int: state = _BotState(out_dir=cfg.out_dir, meeting_id=_meeting_id_from_url(cfg.url), url=cfg.url) - # SIGTERM sets a flag (not an exception) so the Playwright teardown below - # still runs and ``meet_leave`` gets a finalized transcript. + # 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} + for sig in (signal.SIGTERM, signal.SIGINT): + signal.signal(sig, lambda _sig, _frame: stop_flag.__setitem__("stop", True)) - def _on_signal(_sig, _frame): - stop_flag["stop"] = True - - signal.signal(signal.SIGTERM, _on_signal) - signal.signal(signal.SIGINT, _on_signal) - - # 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} + # Realtime resources 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"]: _setup_realtime(rt, cfg.realtime_api_key, state) @@ -497,20 +460,17 @@ def run_bot() -> int: from playwright.sync_api import sync_playwright except ImportError as e: state.set(error=f"playwright not installed: {e}", exited=True) - sys.stderr.write( - "google_meet bot: playwright is not installed. Run " - "`pip install playwright && python -m playwright install chromium`\n") + sys.stderr.write("google_meet bot: playwright is not installed. Run " + "`pip install playwright && python -m playwright install chromium`\n") if rt["bridge"]: rt["bridge"].teardown() return 3 chrome_args = ["--use-fake-ui-for-media-stream", "--disable-blink-features=AutomationControlled"] if not rt["enabled"]: - # Silent fake device — mic content is irrelevant when we're not speaking. - chrome_args.insert(1, "--use-fake-device-for-media-stream") + chrome_args.insert(1, "--use-fake-device-for-media-stream") # silent fake mic elif rt["bridge_info"] and rt["bridge_info"].get("platform") == "linux": - # Playwright's launch() takes no env: set PULSE_SOURCE on our own - # process so the child Chrome inherits the virtual source. + # Playwright's launch() takes no env: set PULSE_SOURCE on ourselves so Chrome inherits it. os.environ["PULSE_SOURCE"] = rt["bridge_info"].get("device_name", "") try: @@ -518,9 +478,8 @@ def run_bot() -> int: browser = pw.chromium.launch(headless=not cfg.headed, args=chrome_args) context_args = { "viewport": {"width": 1280, "height": 800}, - "user_agent": ( - "Mozilla/5.0 (X11; Linux x86_64) AppleWebKit/537.36 " - "(KHTML, like Gecko) Chrome/124.0.0.0 Safari/537.36"), + "user_agent": ("Mozilla/5.0 (X11; Linux x86_64) AppleWebKit/537.36 " + "(KHTML, like Gecko) Chrome/124.0.0.0 Safari/537.36"), "permissions": ["microphone", "camera"]} if cfg.auth_state and Path(cfg.auth_state).is_file(): context_args["storage_state"] = cfg.auth_state @@ -561,14 +520,10 @@ def run_bot() -> int: def _looks_like_human_speaker(speaker: str, bot_guest_name: str) -> bool: - """Whether a caption's speaker is probably a human rather than our own echo. - - 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 - return speaker.strip().lower() not in {"unknown", "you", bot_guest_name.strip().lower()} + """Whether a caption's speaker is probably a human rather than our own echo (Meet attributes + our fake-mic audio to the bot's name; blank/unknown speakers are ambiguous — no barge-in).""" + return bool(speaker and speaker.strip()) and ( + speaker.strip().lower() not in {"unknown", "you", bot_guest_name.strip().lower()}) _DURATION_UNITS = {"h": 3600.0, "m": 60.0, "s": 1.0} diff --git a/plugins/google_meet/node/__init__.py b/plugins/google_meet/node/__init__.py index 91c6677a61..88696a7f94 100644 --- a/plugins/google_meet/node/__init__.py +++ b/plugins/google_meet/node/__init__.py @@ -1,12 +1,10 @@ -"""Remote 'node host' primitive: run the Meet bot on a different machine than the gateway. +"""Remote 'node host': run the Meet bot on a different machine than the gateway. gateway (Linux) ── ws://mac.local:18789 ──▶ node server (Mac) → process_manager → meet_bot -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. - -NodeClient (gateway-side RPC), NodeServer (hosts the bot), NodeRegistry -(approved nodes: name → url+token), protocol (envelope helpers). +Why: Google sign-in + Chrome profile live on the user's laptop; running the bot there reuses +that profile without shipping credentials. NodeClient (gateway RPC), NodeServer (hosts the +bot), NodeRegistry (approved nodes: name → url+token), protocol (envelope helpers). """ from __future__ import annotations diff --git a/plugins/google_meet/node/cli.py b/plugins/google_meet/node/cli.py index 14ab77c523..56e1d222e5 100644 --- a/plugins/google_meet/node/cli.py +++ b/plugins/google_meet/node/cli.py @@ -41,8 +41,7 @@ def _cmd_run(args: argparse.Namespace, reg: NodeRegistry) -> int: 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"[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()) diff --git a/plugins/google_meet/node/client.py b/plugins/google_meet/node/client.py index 93ce8db7d6..5cace01f95 100644 --- a/plugins/google_meet/node/client.py +++ b/plugins/google_meet/node/client.py @@ -1,9 +1,5 @@ -"""Gateway-side RPC client for a remote meet node. - -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. -""" +"""Gateway-side RPC client for a remote meet node: one short-lived sync WebSocket per call, so +non-async tool handlers need no persistent connection. ``websockets`` is imported lazily.""" from __future__ import annotations @@ -29,10 +25,8 @@ class NodeClient: try: from websockets.sync.client import connect # type: ignore except ImportError as exc: - raise RuntimeError( - "NodeClient requires the 'websockets' package. " - "Install it with: pip install websockets" - ) from exc + raise RuntimeError("NodeClient requires the 'websockets' package. " + "Install it with: pip install websockets") from exc req = _proto.make_request(type, self.token, payload) with connect(self.url, open_timeout=self.timeout, close_timeout=self.timeout) as ws: @@ -48,13 +42,8 @@ class NodeClient: raise RuntimeError("response missing payload dict") return payload_out - def start_bot( - self, - url: str, - guest_name: str = "Hermes Agent", - duration: Optional[str] = None, - headed: bool = False, - mode: str = "transcribe") -> Dict[str, Any]: + def start_bot(self, url: str, guest_name: str = "Hermes Agent", duration: Optional[str] = None, + headed: bool = False, mode: str = "transcribe") -> Dict[str, Any]: payload: Dict[str, Any] = {"url": url, "guest_name": guest_name, "headed": bool(headed), "mode": mode} if duration is not None: payload["duration"] = duration diff --git a/plugins/google_meet/node/protocol.py b/plugins/google_meet/node/protocol.py index 5e32643739..623fe6333c 100644 --- a/plugins/google_meet/node/protocol.py +++ b/plugins/google_meet/node/protocol.py @@ -1,14 +1,11 @@ -"""Wire protocol for gateway ↔ node RPC. - -Everything is a JSON object with the same envelope shape: +"""Wire protocol for gateway ↔ node RPC (JSON envelopes). Request: {"type": , "id": , "token": , "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 carry the shared bearer token (``hermes meet node approve`` on the gateway, read off +disk on the server); mismatched tokens are rejected before dispatch. """ from __future__ import annotations @@ -55,11 +52,8 @@ def encode(msg: Dict[str, Any]) -> str: def decode(raw) -> Dict[str, Any]: - """Parse a JSON envelope (object with string ``type`` + ``id``), raising ValueError otherwise. - - Accepts ``str`` or UTF-8 ``bytes``. Token match and payload shape are - checked server-side in :func:`validate_request`. - """ + """Parse a JSON envelope (object with string ``type`` + ``id``) from str/bytes; ValueError otherwise. + Token match and payload shape are checked server-side in :func:`validate_request`.""" if isinstance(raw, (bytes, bytearray)): raw = raw.decode("utf-8") try: diff --git a/plugins/google_meet/node/registry.py b/plugins/google_meet/node/registry.py index 4755f627e4..2a5fb08cbe 100644 --- a/plugins/google_meet/node/registry.py +++ b/plugins/google_meet/node/registry.py @@ -26,8 +26,6 @@ class NodeRegistry: 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, Dict[str, Any]]: """The ``nodes`` map (name → entry); empty when the file is missing or malformed.""" data = read_json(self.path) @@ -37,8 +35,6 @@ class NodeRegistry: 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]]: entry = self._load().get(name) return None if entry is None else {"name": name, **entry} @@ -63,10 +59,8 @@ class NodeRegistry: return [{"name": name, **entry} for name, entry in sorted(self._load().items())] def resolve(self, chrome_node: Optional[str]) -> Optional[Dict[str, Any]]: - """Named node's entry, or — when ``chrome_node`` is falsy — the sole registered node. - - None when the name is unknown or when zero / several nodes are registered (ambiguous). - """ + """Named node's entry, or (``chrome_node`` falsy) the sole registered node; None if unknown + or when zero / several nodes are registered (ambiguous).""" if chrome_node: return self.get(chrome_node) nodes = self.list_all() diff --git a/plugins/google_meet/node/server.py b/plugins/google_meet/node/server.py index 9c1c9e80e3..b38fda7b4a 100644 --- a/plugins/google_meet/node/server.py +++ b/plugins/google_meet/node/server.py @@ -1,15 +1,9 @@ -"""Remote node server — hosts the Meet bot on another machine (e.g. the user's Mac). +"""Remote node server — hosts the Meet bot on another machine (``hermes meet node run``). -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``. - -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 ``. - -``websockets`` is optional and imported lazily inside :meth:`serve`. +WebSocket endpoint accepting token-signed RPC requests dispatched to ``process_manager``. +Token: 32 hex chars minted on first boot, persisted at ``$HERMES_HOME/workspace/meetings/ +node_token.json`` so approved gateways survive restarts; the operator copies it to the gateway +via ``hermes meet node approve ``. ``websockets`` is imported lazily. """ from __future__ import annotations @@ -41,9 +35,7 @@ def _rpc_start_bot(payload: Dict[str, Any], pm) -> Dict[str, Any]: 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". + # The bot-side consumer only exists in realtime mode: ok=True means "enqueued", not "spoken". text = payload.get("text", "") active = pm._read_active() enqueued = False @@ -71,12 +63,8 @@ _RPC = { class NodeServer: """WebSocket server that executes meet bot RPCs locally.""" - def __init__( - self, - host: str = "127.0.0.1", - port: int = 18789, - token_path: Optional[Path] = None, - display_name: str = "hermes-meet-node") -> None: + def __init__(self, host: str = "127.0.0.1", port: int = 18789, token_path: Optional[Path] = None, + display_name: str = "hermes-meet-node") -> None: self.host = host self.port = port self.display_name = display_name @@ -99,18 +87,16 @@ class NodeServer: async def _handle_request(self, msg: Dict[str, Any]) -> Dict[str, Any]: """Validate + dispatch one decoded request; always returns an envelope, never raises. - - 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. - """ + Envelope ``error`` is for auth/protocol failures and pm crashes; pm's own ``ok``/``error`` + results travel inside a normal response payload.""" ok, reason = _proto.validate_request(msg, self.ensure_token()) if not ok: return _proto.make_error(str(msg.get("id") or ""), reason) 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()}} + 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}") @@ -130,10 +116,8 @@ class NodeServer: try: import websockets # type: ignore except ImportError as exc: - raise RuntimeError( - "NodeServer.serve requires the 'websockets' package. " - "Install it with: pip install websockets" - ) from exc + raise RuntimeError("NodeServer.serve requires the 'websockets' package. " + "Install it with: pip install websockets") from exc self.ensure_token() diff --git a/plugins/google_meet/process_manager.py b/plugins/google_meet/process_manager.py index 1625ad4720..a263483926 100644 --- a/plugins/google_meet/process_manager.py +++ b/plugins/google_meet/process_manager.py @@ -1,17 +1,9 @@ """Subprocess lifecycle manager for the google_meet bot. -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. - -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 +One active meeting at a time, recorded in ``$HERMES_HOME/workspace/meetings/.active.json`` +(``pid, meeting_id, out_dir, url, started_at, session_id, log_path, mode``) so tool calls +across turns can find the bot. The bot is a detached subprocess reached via files only +(``/status.json``, ``/transcript.txt``), so the agent loop can't block. """ from __future__ import annotations @@ -43,8 +35,7 @@ def _write_active(data: Dict[str, Any]) -> None: def _pid_alive(pid: int) -> bool: - # Not ``os.kill(pid, 0)``: on Windows that routes through - # GenerateConsoleCtrlEvent and can kill the target (bpo-14484). + # Not ``os.kill(pid, 0)``: on Windows that can kill the target (bpo-14484). from gateway.status import _pid_exists return bool(pid) and _pid_exists(pid) @@ -64,22 +55,12 @@ def _kill(pid: int, sig) -> None: _NO_ACTIVE = {"ok": False, "reason": "no active meeting"} -def start( - url: str, - *, - out_dir: Optional[Path] = None, - headed: bool = False, - auth_state: Optional[str] = None, - guest_name: str = "Hermes Agent", - duration: Optional[str] = None, - session_id: Optional[str] = None, - mode: str = "transcribe", - realtime_model: Optional[str] = None, - realtime_voice: Optional[str] = None, - realtime_instructions: Optional[str] = None, - realtime_api_key: Optional[str] = None) -> Dict[str, Any]: - """Spawn the meet_bot subprocess for *url*, stopping any running bot first - (single-active-meeting semantics). Returns a dict summarizing the bot.""" +def start(url: str, *, out_dir: Optional[Path] = None, headed: bool = False, + auth_state: Optional[str] = None, guest_name: str = "Hermes Agent", duration: Optional[str] = None, + session_id: Optional[str] = None, mode: str = "transcribe", realtime_model: Optional[str] = None, + realtime_voice: Optional[str] = None, realtime_instructions: Optional[str] = None, + realtime_api_key: Optional[str] = None) -> Dict[str, Any]: + """Spawn the meet_bot subprocess for *url*, stopping any running bot first (one active meeting).""" from plugins.google_meet.meet_bot import _is_safe_meet_url, _meeting_id_from_url if not _is_safe_meet_url(url): @@ -92,7 +73,7 @@ def start( out = out_dir or (_root() / meeting_id) out.mkdir(parents=True, exist_ok=True) - # Wipe stale 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. for name in ("transcript.txt", "status.json"): try: (out / name).unlink() @@ -111,10 +92,8 @@ def start( (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). + # Resolve the realtime key at SPAWN time in the parent, where the profile secret scope + # (a contextvar) is installed; the detached child inherits env, not scope. if not realtime_api_key: from agent.secret_scope import get_secret @@ -123,8 +102,7 @@ def start( 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. The child owns the log fd after Popen. + # Detach: stdout/stderr → log file, new session so parent signals don't propagate. with open(log_path, "ab", buffering=0) as log_fh: proc = subprocess.Popen( [sys.executable, "-m", "plugins.google_meet.meet_bot"], @@ -165,19 +143,13 @@ def transcript(last: Optional[int] = None) -> Dict[str, Any]: tp = Path(active.get("out_dir", "")) / "transcript.txt" 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()] - return { - "ok": True, - "meetingId": active.get("meeting_id"), - "lines": all_lines[-last:] if last else all_lines, - "total": len(all_lines), - "path": str(tp)} + return {"ok": True, "meetingId": active.get("meeting_id"), + "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 ``/say_queue.jsonl`` for the bot's speaker thread. - - Refused when no meeting is active or the active bot is in transcribe-only mode. - """ + """Append a ``say`` request to ``/say_queue.jsonl``. + Refused when no meeting is active or the active bot is transcribe-only.""" import uuid text = (text or "").strip() @@ -225,8 +197,5 @@ def stop(*, reason: str = "requested") -> Dict[str, Any]: (_root() / ".active.json").unlink() except FileNotFoundError: pass - return { - "ok": True, - "reason": reason, - "meetingId": active.get("meeting_id"), - "transcriptPath": str(Path(out_dir) / "transcript.txt") if out_dir else None} + return {"ok": True, "reason": reason, "meetingId": active.get("meeting_id"), + "transcriptPath": str(Path(out_dir) / "transcript.txt") if out_dir else None} diff --git a/plugins/google_meet/realtime/__init__.py b/plugins/google_meet/realtime/__init__.py index 37eb16add3..836675d201 100644 --- a/plugins/google_meet/realtime/__init__.py +++ b/plugins/google_meet/realtime/__init__.py @@ -1,9 +1,4 @@ -"""Realtime speech subpackage for the google_meet plugin (v2). - -Provides a thin OpenAI Realtime API client and a file-queue speaker -wrapper so the Meet bot can play synthesized speech through the -virtual audio bridge. -""" +"""Realtime speech: thin OpenAI Realtime client + file-queue speaker for the Meet bot.""" from .openai_client import RealtimeSession, RealtimeSpeaker # noqa: F401 diff --git a/plugins/google_meet/realtime/openai_client.py b/plugins/google_meet/realtime/openai_client.py index e8c9523520..820b3919b4 100644 --- a/plugins/google_meet/realtime/openai_client.py +++ b/plugins/google_meet/realtime/openai_client.py @@ -1,9 +1,7 @@ """OpenAI Realtime API WebSocket client + file-queue speaker. -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. +text → OpenAI Realtime → audio deltas appended as PCM to a file the audio bridge streams into +Chrome's fake mic. One sync WebSocket per session; ``websockets`` is imported lazily. """ from __future__ import annotations @@ -30,20 +28,11 @@ def _decode_audio(b64: str) -> bytes: class RealtimeSession: - """Minimal sync client for the OpenAI Realtime WebSocket API. + """Minimal sync client for the OpenAI Realtime WebSocket API; ``speak`` and ``cancel_response`` + may run on different threads — a lock serializes WebSocket writes.""" - ``speak`` and ``cancel_response`` may be called from different threads; - a lock serializes WebSocket writes. - """ - - def __init__( - self, - api_key: str, - model: str = "gpt-realtime", - voice: str = "alloy", - instructions: str = "", - audio_sink_path: Optional[Path] = None, - sample_rate: int = 24000) -> None: + def __init__(self, api_key: str, model: str = "gpt-realtime", voice: str = "alloy", + instructions: str = "", audio_sink_path: Optional[Path] = None, sample_rate: int = 24000) -> None: self.api_key = api_key self.model = model self.voice = voice @@ -52,8 +41,7 @@ class RealtimeSession: self.sample_rate = sample_rate self._ws: Any = None self._send_lock = threading.Lock() - # Public counters for status reporting. - self.audio_bytes_out: int = 0 + self.audio_bytes_out: int = 0 # public counters for status reporting self.last_audio_out_at: Optional[float] = None def connect(self) -> None: @@ -61,10 +49,8 @@ class RealtimeSession: 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 + 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")] # Newer websockets takes additional_headers=, older extra_headers=. @@ -91,12 +77,8 @@ class RealtimeSession: self._ws = None def speak(self, text: str, timeout: float = 30.0) -> dict: - """Send ``text`` and append the audio response to ``audio_sink_path``. - - 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. - """ + """Send ``text`` and append the audio response to ``audio_sink_path`` (opened 'ab' per call + so a streaming reader can consume it). Frames other than audio deltas/terminal/error are ignored.""" if self._ws is None: raise RuntimeError("RealtimeSession.connect() must be called first") @@ -153,9 +135,7 @@ class RealtimeSession: 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. - """ + Non-dict / unparseable frames are skipped; TimeoutError past the deadline.""" assert self._ws is not None while True: remaining = deadline - time.monotonic() @@ -176,15 +156,10 @@ class RealtimeSession: class RealtimeSpeaker: - """File-based JSONL queue wrapper around :class:`RealtimeSession`. + """JSONL queue (``{"id", "text"}`` per line) wrapper around :class:`RealtimeSession`; processed + lines are appended to ``processed_path`` (if set) and removed from the queue.""" - Each queue line is ``{"id": "", "text": "..."}``. Processed lines - are appended to ``processed_path`` (if set) and removed from the queue. - """ - - def __init__( - self, session: RealtimeSession, queue_path: Path, processed_path: Optional[Path] = None - ) -> None: + def __init__(self, session: RealtimeSession, queue_path: Path, processed_path: Optional[Path] = None) -> None: self.session = session self.queue_path = Path(queue_path) self.processed_path = Path(processed_path) if processed_path else None @@ -205,8 +180,7 @@ class RealtimeSpeaker: return out def _rewrite_queue(self, remaining: list[dict]) -> None: - # Always keep the file (empty when drained): consumers may watch its - # mtime, and delete-then-recreate is a race. + # Always keep the file (empty when drained): consumers may watch its mtime. body = "".join(json.dumps(e) + "\n" for e in remaining) self.queue_path.write_text(body, encoding="utf-8") @@ -236,8 +210,7 @@ class RealtimeSpeaker: result = {"ok": True, "bytes_written": 0, "duration_ms": 0.0} self._append_processed(head, result) - # Re-read from disk (new entries may have arrived), then drop the - # head — by position when it's still first, else by id. + # Re-read (new entries may have arrived), then drop the head by position or id. latest = self._read_queue() if latest and latest[0].get("id") == head.get("id"): self._rewrite_queue(latest[1:]) diff --git a/plugins/google_meet/tools.py b/plugins/google_meet/tools.py index 61df319d22..a2d97467fc 100644 --- a/plugins/google_meet/tools.py +++ b/plugins/google_meet/tools.py @@ -17,18 +17,11 @@ from plugins.google_meet import process_manager as pm def check_meet_requirements() -> bool: """True when the plugin can run LOCALLY: Linux/macOS + importable ``playwright``. - - Remote-node operation only needs ``websockets`` on the gateway side; the - handlers relax this gate themselves when a node is addressed. - """ + Remote-node operation only needs ``websockets``; handlers relax this gate when a node is addressed.""" + import importlib.util import platform as _p - if _p.system().lower() not in {"linux", "darwin"}: - return False - try: - import playwright # noqa: F401 - except ImportError: - return False - return True + return (_p.system().lower() in {"linux", "darwin"} + and importlib.util.find_spec("playwright") is not None) def resolve_node(node: str): @@ -42,10 +35,6 @@ def resolve_node(node: str): return NodeClient(url=entry["url"], token=entry["token"]), entry.get("name") -# --------------------------------------------------------------------------- -# Schemas -# --------------------------------------------------------------------------- - _NODE_PROP = {"type": "string"} MEET_JOIN_SCHEMA: Dict[str, Any] = { @@ -168,10 +157,6 @@ MEET_SAY_SCHEMA: Dict[str, Any] = { } -# --------------------------------------------------------------------------- -# Handlers -# --------------------------------------------------------------------------- - def _json(obj: Any) -> str: return json.dumps(obj, ensure_ascii=False) @@ -187,10 +172,8 @@ def _dispatch(node: Optional[str], op: str, remote, local) -> str: 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" - ) + return _err(f"no registered meet node matches {node!r} — " + "run `hermes meet node approve ` first") try: res = remote(client) except Exception as e: @@ -207,23 +190,16 @@ def handle_meet_join(args: Dict[str, Any], **_kw) -> str: return _err(f"mode must be 'transcribe' or 'realtime' (got {mode!r})") common: Dict[str, Any] = dict( - url=url, - guest_name=str(args.get("guest_name") or "Hermes Agent"), + 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, - ) + headed=bool(args.get("headed", False)), mode=mode) 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 {"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) @@ -240,17 +216,13 @@ def handle_meet_transcript(args: Dict[str, Any], **_kw) -> str: 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), - ) + 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: - return _dispatch( - args.get("node"), "stop", - lambda c: c.stop(), lambda: pm.stop(reason="agent called meet_leave"), - ) + 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: diff --git a/plugins/plugin_loader.py b/plugins/plugin_loader.py index 9526476064..6ef572bf63 100644 --- a/plugins/plugin_loader.py +++ b/plugins/plugin_loader.py @@ -1,10 +1,7 @@ -"""Shared directory-plugin loader for the ``plugins///`` discovery packages. - -Used by ``plugins.cron_providers`` and ``plugins.context_engine``: import a plugin -``__init__.py`` by path (with its sibling ``*.py`` pre-registered so relative -imports work), then extract the provider instance via ``register(ctx)`` or an -ABC-subclass fallback. -""" +"""Shared directory-plugin loader for ``plugins///`` discovery packages +(cron_providers, context_engine, memory): import ``__init__.py`` by path with siblings +pre-registered so relative imports work, then extract the provider via ``register(ctx)`` +or an ABC-subclass fallback.""" from __future__ import annotations @@ -41,11 +38,8 @@ def iter_plugin_dirs(root: Path) -> List[Path]: """Sorted child dirs of *root* that have an ``__init__.py`` (skips ``_``/``.`` names).""" if not root.is_dir(): return [] - return [ - child for child in sorted(root.iterdir()) - if child.is_dir() - and not child.name.startswith(("_", ".")) - and (child / "__init__.py").exists()] + return [child for child in sorted(root.iterdir()) + if child.is_dir() and not child.name.startswith(("_", ".")) and (child / "__init__.py").exists()] def read_plugin_description(plugin_dir: Path) -> str: @@ -75,10 +69,7 @@ def _new_module(name: str, file: Path, search_locations: Optional[List[str]] = N def _exec(mod: Any, logger: Optional[logging.Logger] = None) -> bool: """Execute a module made by ``_new_module`` (None -> False); False (and debug-log) if it raised. - - The sys.modules entry is left in place on failure; callers that need a clean - retry (the main plugin module) pop it themselves. - """ + The sys.modules entry stays on failure; callers needing a clean retry pop it themselves.""" if mod is None: return False try: @@ -90,19 +81,13 @@ def _exec(mod: Any, logger: Optional[logging.Logger] = None) -> bool: return False -def load_plugin_module( - module_name: str, - plugin_dir: Path, - *, - parents: Tuple[str, ...], - logger: logging.Logger, - synthetic_namespace: Optional[str] = None) -> Optional[Any]: +def load_plugin_module(module_name: str, plugin_dir: Path, *, parents: Tuple[str, ...], + logger: logging.Logger, synthetic_namespace: Optional[str] = None) -> Optional[Any]: """Import ``plugin_dir/__init__.py`` as *module_name* (reusing sys.modules when loaded). - Order matters: parent packages first (relative imports need them), then sibling - ``*.py`` as ``module_name.`` (so ``from ._x import Y`` resolves), then the - module itself. Finally the child is bound onto its parent and the siblings onto the - module — the shape normal imports produce, which dotted imports and monkeypatch rely on. + Order matters: parents first (relative imports need them), then siblings as ``module_name.`` + (so ``from ._x import Y`` resolves), then the module. Finally child is bound onto parent and + siblings onto module — the shape normal imports produce, which monkeypatch relies on. """ init_file = plugin_dir / "__init__.py" if not init_file.exists(): @@ -148,8 +133,8 @@ def load_plugin_module( class NoopPluginContext: - """Base for the fake ``register(ctx)`` contexts: every registration is a no-op except - the one the subclass overrides to capture its provider.""" + """Base for fake ``register(ctx)`` contexts: every registration is a no-op except the one the + subclass overrides to capture its provider.""" def _noop(self, *args, **kwargs): pass @@ -157,14 +142,8 @@ class NoopPluginContext: register_tool = register_hook = register_cli_command = register_memory_provider = _noop -def instance_from_module( - mod: Any, - *, - collector: Any, - collected_attr: str, - base_cls: type, - name: str, - logger: logging.Logger) -> Optional[Any]: +def instance_from_module(mod: Any, *, collector: Any, collected_attr: str, base_cls: type, name: str, + logger: logging.Logger) -> Optional[Any]: """Extract the provider instance: ``register(ctx)`` first, then any ``base_cls`` subclass.""" if hasattr(mod, "register"): try: @@ -185,14 +164,8 @@ def instance_from_module( return None -def load_named( - name: str, - plugin_dir: Path, - load_from_dir: Callable[[Path], Optional[Any]], - *, - kind: str, - noun: str, - logger: logging.Logger) -> Optional[Any]: +def load_named(name: str, plugin_dir: Path, load_from_dir: Callable[[Path], Optional[Any]], *, kind: str, + noun: str, logger: logging.Logger) -> Optional[Any]: """Shared body of ``load_(name)``: load from *plugin_dir*, warn + None on failure.""" try: instance = load_from_dir(plugin_dir) diff --git a/plugins/plugin_storage.py b/plugins/plugin_storage.py index baba980927..9629038500 100644 --- a/plugins/plugin_storage.py +++ b/plugins/plugin_storage.py @@ -1,17 +1,9 @@ -"""Per-plugin persistent storage convention. +"""Per-plugin persistent storage: ``/plugin-data//``. -Plugins must NOT park state inside ``/plugins//`` — that is the -install dir, which ``hermes plugins remove`` deletes and ``update`` git-pulls into. The -sanctioned root is ``/plugin-data//``: user-owned, untouched by -install/update/remove, inspectable in one place. Secrets are deliberately NOT part of -this convention — credential reads go through ``agent.secret_scope`` / ``.env``. - -Usage:: - - from plugins.plugin_storage import plugin_data_dir, plugin_db - - state_file = plugin_data_dir("my-plugin") / "state.json" - conn = plugin_db("my-plugin") # /data.db +Plugins must NOT park state in ``/plugins//`` (the install dir, deleted by +``remove`` and git-pulled by ``update``). Secrets are deliberately NOT part of this convention — +credential reads go through ``agent.secret_scope`` / ``.env``. +Usage: ``plugin_data_dir("my-plugin") / "state.json"``; ``plugin_db("my-plugin")`` → ``data.db``. """ from __future__ import annotations @@ -22,8 +14,7 @@ from pathlib import Path __all__ = ["plugin_data_dir", "plugin_db"] -# Mirrors the plugin-name shape `hermes plugins install` accepts; anything else could -# escape the data root via separators or traversal. +# Mirrors the plugin-name shape `hermes plugins install` accepts (no separators/traversal). _NAME_RE = re.compile(r"^[a-zA-Z0-9][a-zA-Z0-9._-]{0,63}$") @@ -34,11 +25,8 @@ def _validate_name(name: str) -> str: def plugin_data_dir(name: str) -> Path: - """Return (and create) ``/plugin-data//``. - - Resolves ``get_hermes_home()`` on every call so it follows the active profile — - don't cache the result across profile switches. - """ + """Return (and create) ``/plugin-data//``; resolves ``get_hermes_home()`` on + every call so it follows the active profile — don't cache across profile switches.""" from hermes_constants import get_hermes_home root = get_hermes_home() / "plugin-data" / _validate_name(name) @@ -47,11 +35,8 @@ def plugin_data_dir(name: str) -> Path: def plugin_db(name: str, filename: str = "data.db") -> sqlite3.Connection: - """Open ``/`` (created on first use). - - WAL so a dashboard reader and an agent-tool writer coexist; ``check_same_thread=False`` - matches the multi-threaded FastAPI/tool environment — the caller owns transaction discipline. - """ + """Open ``/``. WAL so a dashboard reader and a tool writer coexist; + ``check_same_thread=False`` for the threaded FastAPI/tool env — caller owns transactions.""" if Path(filename).name != filename or not filename: raise ValueError(f"invalid plugin db filename: {filename!r}") diff --git a/plugins/plugin_utils.py b/plugins/plugin_utils.py index 36ea489527..c3480c06b1 100644 --- a/plugins/plugin_utils.py +++ b/plugins/plugin_utils.py @@ -1,15 +1,8 @@ -"""Shared concurrency helpers for plugin authors. +"""Thread-safe lazy singletons for plugin authors (stdlib-only). -The common plugin footgun is the lazy process-wide singleton (``if _client is None: -_client = Expensive()``): two threads both pass the guard, both build, and the second -write leaks the first's connections/threads. Multi-threaded agent sessions (delegated -tool calls, background workers) make this reachable, so use these instead of hand-rolling -double-checked locking: - -* :func:`lazy_singleton` — decorator for the zero-arg accessor case. -* :class:`SingletonSlot` — manual slot when the instance depends on a config/key argument. - -Both are stdlib-only (``threading``) so any plugin can import them cheaply. +The ``if _client is None: _client = Expensive()`` footgun: two threads pass the guard, both +build, the second leaks the first's connections. :func:`lazy_singleton` decorates a zero-arg +accessor; :class:`SingletonSlot` is the manual slot when the instance depends on an argument. """ from __future__ import annotations @@ -24,18 +17,9 @@ T = TypeVar("T") class SingletonSlot(Generic[T]): - """Thread-safe lazy slot for accessors that take a build argument. - - Caches the first successfully-built instance and ignores the argument afterwards - ("first config wins", the semantics most plugins rely on). The factory runs at most - once under concurrent first calls; if it raises, nothing is cached and the next call - retries. Example:: - - _slot: SingletonSlot[Honcho] = SingletonSlot() - - def get_honcho_client(config=None): - return _slot.get(lambda: Honcho(**resolve(config))) - """ + """Thread-safe lazy slot: caches the first successfully-built instance ("first config wins"). + The factory runs at most once under concurrent first calls; if it raises, nothing is cached + and the next call retries. ``_slot.get(lambda: Honcho(**resolve(config)))``.""" __slots__ = ("_lock", "_value", "_set") @@ -68,15 +52,8 @@ class SingletonSlot(Generic[T]): def lazy_singleton(factory: Callable[[], T]) -> Callable[[], T]: - """Wrap a zero-argument factory into a thread-safe lazy singleton accessor. - - The factory runs exactly once even under concurrent first calls (if it raises, the - next call retries). A ``.reset()`` attribute drops the instance for tests/teardown:: - - @lazy_singleton - def get_client(): - return ExpensiveClient(load_config()) - """ + """Wrap a zero-argument factory into a thread-safe lazy singleton accessor (factory runs once + even under concurrent first calls; on raise the next call retries). ``.reset()`` drops it.""" slot: SingletonSlot[T] = SingletonSlot() @functools.wraps(factory)