refactor(plugins): drop dead disk_cleanup.deep + unused constants, compact docstrings/comments, tighten guards in google_meet/chronos/loader

This commit is contained in:
Teknium
2026-09-02 21:34:56 -07:00
parent fc46eabcf5
commit 8c1ba7b160
25 changed files with 381 additions and 863 deletions
+13 -24
View File
@@ -1,10 +1,6 @@
"""Context engine plugin discovery.
Scans ``plugins/context_engine/<name>/`` 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/<name>/`` → ``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)
+10 -22
View File
@@ -1,11 +1,7 @@
"""Cron scheduler provider plugin discovery.
Scans bundled ``plugins/cron_providers/<name>/`` then user ``$HERMES_HOME/plugins/<name>/``
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/<name>/`` then user
``$HERMES_HOME/plugins/<name>/`` (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
+26 -43
View File
@@ -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"]
+11 -22
View File
@@ -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]:
+15 -30
View File
@@ -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)
+19 -48
View File
@@ -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 <path>"
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.")
+48 -129
View File
@@ -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/<name>`` — state and audit log deliberately
live outside ``$HERMES_HOME/logs/``."""
"""``$HERMES_HOME/disk-cleanup/<name>`` — 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
+10 -20
View File
@@ -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)
+2 -5
View File
@@ -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")
+11 -29
View File
@@ -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=<source_name>`` 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=<source_name>`` 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
+33 -65
View File
@@ -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:
+59 -104
View File
@@ -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}
+4 -6
View File
@@ -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
+1 -2
View File
@@ -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 <name> ws://<host>:{args.port} {token}")
try:
asyncio.run(server.serve())
+6 -17
View File
@@ -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
+5 -11
View File
@@ -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": <str>, "id": <str>, "token": <str>, "payload": <dict>}
Response: {"type": "response", "id": <req-id>, "payload": <dict>}
Error: {"type": "error", "id": <req-id>, "error": <str>}
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:
+2 -8
View File
@@ -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()
+14 -30
View File
@@ -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 <name> <url> <token>``.
``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 <name> <url> <token>``. ``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()
+21 -52
View File
@@ -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"}
<meeting-id>/status.json live bot state (written by the bot)
<meeting-id>/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
(``<meeting-id>/status.json``, ``<meeting-id>/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 ``<out_dir>/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 ``<out_dir>/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}
+1 -6
View File
@@ -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
+17 -44
View File
@@ -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": "<uuid>", "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:])
+16 -44
View File
@@ -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 <name> <url> <token>` first"
)
return _err(f"no registered meet node matches {node!r} — "
"run `hermes meet node approve <name> <url> <token>` 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:
+18 -45
View File
@@ -1,10 +1,7 @@
"""Shared directory-plugin loader for the ``plugins/<kind>/<name>/`` 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/<kind>/<name>/`` 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.<stem>`` (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.<stem>``
(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_<kind>(name)``: load from *plugin_dir*, warn + None on failure."""
try:
instance = load_from_dir(plugin_dir)
+10 -25
View File
@@ -1,17 +1,9 @@
"""Per-plugin persistent storage convention.
"""Per-plugin persistent storage: ``<hermes home>/plugin-data/<name>/``.
Plugins must NOT park state inside ``<hermes home>/plugins/<name>/`` — that is the
install dir, which ``hermes plugins remove`` deletes and ``update`` git-pulls into. The
sanctioned root is ``<hermes home>/plugin-data/<name>/``: 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 dir>/data.db
Plugins must NOT park state in ``<hermes home>/plugins/<name>/`` (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) ``<hermes home>/plugin-data/<name>/``.
Resolves ``get_hermes_home()`` on every call so it follows the active profile —
don't cache the result across profile switches.
"""
"""Return (and create) ``<hermes home>/plugin-data/<name>/``; 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 ``<data dir>/<filename>`` (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 ``<data dir>/<filename>``. 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}")
+9 -32
View File
@@ -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)