refactor(plugins): schema builder for meet tools, SimpleNamespace bot config, lru-cached cron guard, dispatch/suppress collapses across loader, chronos, node, realtime

This commit is contained in:
Teknium
2026-09-03 00:00:18 -07:00
parent 5f70de0d59
commit aec3f4bf81
12 changed files with 132 additions and 203 deletions

View File

@@ -52,7 +52,6 @@ class _EngineCollector(_loader.NoopPluginContext):
def __init__(self, engine_name: str = ""):
self.engine = None
self._engine_name = engine_name or "context_engine"
self._registered_commands: list[str] = []
def register_context_engine(self, engine):
self.engine = engine
@@ -80,7 +79,6 @@ class _EngineCollector(_loader.NoopPluginContext):
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()}
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)

View File

@@ -58,23 +58,23 @@ class ChronosCronScheduler(CronScheduler):
self._client = NasCronClient(_cfg("cron", "chronos", "portal_url"))
return self._client
def _reconcile_logged(self, log, what: str) -> None:
try:
self.reconcile()
except Exception as e:
log("Chronos %s reconcile failed: %s", what, e)
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 can't prove what an interrupted process did: classify unknown, never requeue.
self.recover_interrupted()
try:
self.reconcile()
except Exception as e:
logger.warning("Chronos start() reconcile failed: %s", e)
self._reconcile_logged(logger.warning, "start()")
def stop(self) -> None:
pass
def on_jobs_changed(self) -> None:
try:
self.reconcile()
except Exception as e:
logger.debug("Chronos on_jobs_changed reconcile failed: %s", e)
self._reconcile_logged(logger.debug, "on_jobs_changed")
def register_job(self, job: Dict[str, Any]) -> None:
"""Arm the first one-shot for a new job; may raise so creation can report it."""

View File

@@ -64,9 +64,7 @@ def _on_post_tool_call(tool_name: str = "", args: Optional[Dict[str, Any]] = Non
p = Path(path_str).expanduser()
except Exception:
continue
if not p.exists():
continue
category = dg.guess_category(p)
category = dg.guess_category(p) if p.exists() else None
if category is not None and dg.track(str(p), category, silent=True) and category == "test":
with _lock:
_recent_test_tracks.setdefault(task_id or session_id or "default", set()).add(str(p))

View File

@@ -9,6 +9,7 @@ Scope: strictly HERMES_HOME and /tmp/hermes-*; never ~/.hermes/logs/ or system d
from __future__ import annotations
import contextlib
import functools
import json
import logging
import shutil
@@ -98,22 +99,18 @@ _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).
_PROTECTED_CRON_PATHS: set[str] = set()
@functools.lru_cache(maxsize=1) # built lazily so HERMES_HOME resolves once
def _protected_cron_paths() -> frozenset:
"""Defense-in-depth for quick(): EXACT cron control-plane paths (``cron/``, ``output/`` root,
``jobs.json``, ``.tick.lock``) never deleted regardless of stored category (stale tracked.json).
Never widen to everything under ``cron/output/``: run artifacts there are disposable; only
wholesale deletion of ``output/`` is fatal."""
return frozenset(str(x) for parent in ("cron", "cronjobs") for base in (get_hermes_home() / parent,)
for x in (base, base / "output", base / "jobs.json", base / ".tick.lock"))
def _is_protected_cron_path(p: Path) -> bool:
"""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"):
base = hermes_home / parent
_PROTECTED_CRON_PATHS.update(
str(x) for x in (base, base / "output", base / "jobs.json", base / ".tick.lock"))
return str(p.resolve()) in _PROTECTED_CRON_PATHS
return str(p.resolve()) in _protected_cron_paths()
def fmt_size(n: float) -> str:

View File

@@ -42,12 +42,11 @@ class AudioBridge:
def setup(self) -> dict:
"""Provision the device; raises RuntimeError on unsupported platforms or missing tools."""
system = platform.system()
if system == "Linux":
return self._setup_linux()
if system == "Darwin":
return self._setup_darwin()
raise RuntimeError("windows not supported in v2" if system == "Windows"
else f"unsupported platform: {system}")
impl = {"Linux": self._setup_linux, "Darwin": self._setup_darwin}.get(system)
if impl is None:
raise RuntimeError("windows not supported in v2" if system == "Windows"
else f"unsupported platform: {system}")
return impl()
def teardown(self) -> None:
"""Release the virtual audio device. Idempotent; never raises."""

View File

@@ -70,20 +70,19 @@ def register_cli(subparser: argparse.ArgumentParser) -> None:
_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(a.realtime), assume_yes=bool(a.yes)),
"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)),
mode=a.mode, node=a.node),
"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)),
"say": lambda a: _cmd_say(text=a.text, node=a.node),
"stop": lambda a: _print_result(pm.stop(reason="hermes meet stop")),
"node": node_command} # node subparsers are required=True, so a sub-command is always present
def meet_command(args: argparse.Namespace) -> int:
sub = getattr(args, "meet_command", None)
sub = args.meet_command
if not sub:
print("usage: hermes meet {setup,auth,join,status,transcript,say,stop,node}")
return 2
@@ -100,8 +99,7 @@ def _cmd_setup() -> int:
system_ok = system in {"Linux", "Darwin"}
print(f" platform : {system} [{'ok' if system_ok else 'unsupported'}]")
pw_ok = importlib.util.find_spec("playwright") is not None
print(" playwright : installed" if pw_ok
else " playwright : NOT installed — run: pip install playwright")
print(" playwright : " + ("installed" if pw_ok else "NOT installed — run: pip install playwright"))
chromium_ok, chromium_msg = False, "unknown"
if pw_ok:
try:
@@ -114,8 +112,7 @@ def _cmd_setup() -> int:
chromium_msg = f"probe failed: {e}"
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
print("ready. Join a meeting: hermes meet join https://meet.google.com/abc-defg-hij" if all_ok
@@ -182,7 +179,7 @@ def _cmd_install(*, realtime: bool, assume_yes: bool) -> int:
stdin=subprocess.DEVNULL)
except Exception:
have_bh = False
needs = ([] if have_bh else ["blackhole-2ch"]) + ([] if shutil.which("ffmpeg") else ["ffmpeg"])
needs = [pkg for pkg, have in (("blackhole-2ch", have_bh), ("ffmpeg", shutil.which("ffmpeg"))) if not have]
if not needs:
print(" BlackHole and ffmpeg already installed.")
elif not shutil.which("brew"):
@@ -214,8 +211,7 @@ def _cmd_auth() -> int:
with sync_playwright() as pw:
browser = pw.chromium.launch(headless=False)
context = browser.new_context()
page = context.new_page()
page.goto("https://accounts.google.com/", wait_until="domcontentloaded")
context.new_page().goto("https://accounts.google.com/", wait_until="domcontentloaded")
with contextlib.suppress(EOFError):
input("press Enter after you've signed in ... ")
context.storage_state(path=str(path))

View File

@@ -18,8 +18,8 @@ import subprocess
import sys
import threading
import time
from dataclasses import dataclass
from pathlib import Path
from types import SimpleNamespace
from typing import Optional
from plugins.google_meet._jsonfile import write_json_atomic
@@ -255,8 +255,8 @@ def _start_realtime_speaker(rt: dict, cfg: "_BotConfig", stop_flag: dict, state:
state.set(error=f"realtime connect failed: {e}")
return
rt["session"] = session
speaker = RealtimeSpeaker(
session=session, queue_path=queue_path, processed_path=cfg.out_dir / "say_processed.jsonl")
speaker = RealtimeSpeaker(session=session, queue_path=queue_path,
processed_path=cfg.out_dir / "say_processed.jsonl")
def _speaker_loop():
try:
@@ -291,9 +291,8 @@ def _setup_realtime(rt: dict, api_key: str, state: _BotState) -> None:
return
try:
from plugins.google_meet.audio_bridge import AudioBridge
bridge = AudioBridge()
rt["bridge_info"] = bridge.setup()
rt["bridge"] = bridge
rt["bridge"] = AudioBridge()
rt["bridge_info"] = rt["bridge"].setup()
state.set(realtime=True, realtime_device=rt["bridge_info"].get("device_name"))
except Exception as e:
state.set(error=f"audio bridge setup failed: {e} — falling back to transcribe")
@@ -310,22 +309,7 @@ def _teardown_realtime(rt: dict) -> None:
_quiet(getattr(rt[key], method), **kw)
@dataclass
class _BotConfig:
"""Everything the bot reads from ``HERMES_MEET_*`` env vars."""
url: str
out_dir: Optional[Path]
headed: bool
auth_state: str
guest_name: str
duration_s: Optional[float]
realtime: bool
realtime_api_key: str
realtime_model: str
realtime_voice: str
realtime_instructions: str
lobby_timeout: float
_BotConfig = SimpleNamespace # everything the bot reads from ``HERMES_MEET_*`` env vars
def _config_from_env() -> _BotConfig:

View File

@@ -9,6 +9,7 @@ via ``hermes meet node approve <name> <url> <token>``. ``websockets`` is importe
from __future__ import annotations
import asyncio
import contextlib
import json
import secrets
import time
@@ -40,14 +41,12 @@ def _rpc_say(payload: Dict[str, Any], pm) -> Dict[str, Any]:
active = pm._read_active()
enqueued = False
if active and active.get("out_dir"):
try:
with contextlib.suppress(OSError):
queue = Path(active["out_dir"]) / "say_queue.jsonl"
queue.parent.mkdir(parents=True, exist_ok=True)
with queue.open("a", encoding="utf-8") as fh:
fh.write(json.dumps({"text": text, "ts": time.time()}) + "\n")
enqueued = True
except OSError:
pass
return {"ok": True, "enqueued": enqueued, "text": text}

View File

@@ -42,11 +42,6 @@ def _pid_alive(pid: int) -> bool:
return bool(pid) and _pid_exists(pid)
def _active_pid() -> int:
active = _read_active()
return int(active.get("pid", 0)) if active else 0
def _kill(pid: int, sig) -> None:
with contextlib.suppress(ProcessLookupError):
os.kill(pid, sig)
@@ -64,7 +59,7 @@ def start(url: str, *, out_dir: Optional[Path] = None, headed: bool = False,
from plugins.google_meet.meet_bot import _is_safe_meet_url, _meeting_id_from_url
if not _is_safe_meet_url(url):
return {"ok": False, "error": "refusing: only https://meet.google.com/ URLs are allowed. got: " + repr(url)}
if _pid_alive(_active_pid()):
if _pid_alive(int((_read_active() or {}).get("pid", 0))):
stop(reason="replaced by new meet_join")
meeting_id = _meeting_id_from_url(url)
out = out_dir or (_root() / meeting_id)
@@ -163,7 +158,7 @@ def stop(*, reason: str = "requested") -> Dict[str, Any]:
if not _pid_alive(pid):
break
time.sleep(0.5)
if _pid_alive(pid):
else:
_kill(pid, signal.SIGKILL) # windows-footgun: ok — POSIX-only plugin (google_meet registers no-op on Windows; see __init__.py)
(_root() / ".active.json").unlink(missing_ok=True)
return {"ok": True, "reason": reason, "meetingId": active.get("meeting_id"),

View File

@@ -59,14 +59,9 @@ class RealtimeSession:
self._ws = connect(url, additional_headers=headers)
except TypeError:
self._ws = connect(url, extra_headers=headers)
self._send_json({
"type": "session.update",
"session": {
"voice": self.voice,
"instructions": self.instructions,
"modalities": ["audio", "text"],
"output_audio_format": "pcm16",
"input_audio_format": "pcm16"}})
self._send_json({"type": "session.update", "session": {
"voice": self.voice, "instructions": self.instructions, "modalities": ["audio", "text"],
"output_audio_format": "pcm16", "input_audio_format": "pcm16"}})
def close(self) -> None:
if self._ws is not None:
@@ -80,10 +75,8 @@ class RealtimeSession:
if self._ws is None:
raise RuntimeError("RealtimeSession.connect() must be called first")
start = time.monotonic()
self._send_json({
"type": "conversation.item.create",
"item": {"type": "message", "role": "user", "content": [{"type": "input_text", "text": text}]},
})
self._send_json({"type": "conversation.item.create", "item": {
"type": "message", "role": "user", "content": [{"type": "input_text", "text": text}]}})
self._send_json({"type": "response.create", "response": {"modalities": ["audio"]}})
bytes_written = 0
with contextlib.ExitStack() as stack:
@@ -98,14 +91,14 @@ class RealtimeSession:
ftype = frame.get("type")
if ftype == "error":
raise RuntimeError(f"realtime error: {frame.get('error') or frame}")
if ftype == "response.audio.delta" and sink_fp is not None:
chunk = _decode_audio(frame.get("delta") or frame.get("audio") or "")
if chunk:
sink_fp.write(chunk)
sink_fp.flush()
bytes_written += len(chunk)
self.audio_bytes_out += len(chunk)
self.last_audio_out_at = time.time()
chunk = _decode_audio(frame.get("delta") or frame.get("audio") or "") if (
ftype == "response.audio.delta" and sink_fp is not None) else b""
if chunk:
sink_fp.write(chunk)
sink_fp.flush()
bytes_written += len(chunk)
self.audio_bytes_out += len(chunk)
self.last_audio_out_at = time.time()
return {"ok": True, "bytes_written": bytes_written, "duration_ms": (time.monotonic() - start) * 1000.0}
def cancel_response(self) -> bool:
@@ -137,12 +130,10 @@ class RealtimeSession:
raw = self._ws.recv()
if raw is None:
return None
try:
with contextlib.suppress(TypeError, ValueError):
frame = json.loads(raw) if isinstance(raw, (str, bytes, bytearray)) else raw
except (TypeError, ValueError):
continue
if isinstance(frame, dict):
return frame
if isinstance(frame, dict):
return frame
class RealtimeSpeaker:
@@ -160,13 +151,11 @@ class RealtimeSpeaker:
return []
out: list[dict] = []
for line in self.queue_path.read_text(encoding="utf-8").splitlines():
try:
with contextlib.suppress(ValueError):
entry = json.loads(line) if line.strip() else None
except ValueError:
continue
if isinstance(entry, dict):
entry.setdefault("id", str(uuid.uuid4()))
out.append(entry)
if isinstance(entry, dict):
entry.setdefault("id", str(uuid.uuid4()))
out.append(entry)
return out
def _rewrite_queue(self, remaining: list[dict]) -> None:
@@ -200,7 +189,5 @@ class RealtimeSpeaker:
self._append_processed(head, result)
# 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:])
else:
self._rewrite_queue([e for e in latest if e.get("id") != head.get("id")])
self._rewrite_queue(latest[1:] if latest and latest[0].get("id") == head.get("id")
else [e for e in latest if e.get("id") != head.get("id")])

View File

@@ -36,89 +36,73 @@ def resolve_node(node: str):
_NODE_PROP = {"type": "string"}
MEET_JOIN_SCHEMA: Dict[str, Any] = {
"name": "meet_join",
"description": (
"Join a Google Meet call and start scraping live captions into a transcript file. Only "
"meet.google.com URLs are accepted; no calendar scanning, no auto-dial. Spawns a headless "
"Chromium subprocess that runs in parallel with the agent loop — returns immediately. Poll "
"with meet_status and read captions with meet_transcript. Reminder to the agent: you "
"should announce yourself in the meeting (there is no automatic consent announcement)."),
"parameters": {
"type": "object",
"properties": {
"url": {"type": "string", "description": "Full https://meet.google.com/... URL. Required."},
"mode": {
"type": "string", "enum": ["transcribe", "realtime"],
"description": (
"transcribe (default): listen-only, scrape captions. "
"realtime: also enable agent speech via meet_say "
"(requires OpenAI Realtime key + platform audio bridge).")},
"guest_name": {
"type": "string",
"description": "Display name to use when joining as guest. Defaults to 'Hermes Agent'."},
"duration": {
"type": "string",
"description": ("Optional max duration before auto-leave (e.g. '30m', "
"'2h', '90s'). Omit to stay until meet_leave is called.")},
"headed": {
"type": "boolean",
def _str(description: str) -> Dict[str, Any]:
return {"type": "string", "description": description}
def _schema(name: str, description: str, properties: Dict[str, Any], required=None) -> Dict[str, Any]:
params: Dict[str, Any] = {"type": "object", "properties": properties}
if required:
params["required"] = required
params["additionalProperties"] = False
return {"name": name, "description": description, "parameters": params}
MEET_JOIN_SCHEMA = _schema(
"meet_join",
"Join a Google Meet call and start scraping live captions into a transcript file. Only "
"meet.google.com URLs are accepted; no calendar scanning, no auto-dial. Spawns a headless "
"Chromium subprocess that runs in parallel with the agent loop — returns immediately. Poll "
"with meet_status and read captions with meet_transcript. Reminder to the agent: you "
"should announce yourself in the meeting (there is no automatic consent announcement).",
{"url": _str("Full https://meet.google.com/... URL. Required."),
"mode": {"type": "string", "enum": ["transcribe", "realtime"],
"description": ("transcribe (default): listen-only, scrape captions. "
"realtime: also enable agent speech via meet_say "
"(requires OpenAI Realtime key + platform audio bridge).")},
"guest_name": _str("Display name to use when joining as guest. Defaults to 'Hermes Agent'."),
"duration": _str("Optional max duration before auto-leave (e.g. '30m', "
"'2h', '90s'). Omit to stay until meet_leave is called."),
"headed": {"type": "boolean",
"description": "Run Chromium headed instead of headless (debug only). Default false."},
"node": {
"type": "string",
"description": (
"Name of a registered remote node to run the bot on (useful when the gateway "
"runs on a headless Linux box but the user's Chrome with a signed-in Google "
"profile lives on their Mac). Pass 'auto' to use the single registered node. "
"Default: run locally. Nodes are approved via `hermes meet node approve`.")}},
"required": ["url"],
"additionalProperties": False}}
"node": _str("Name of a registered remote node to run the bot on (useful when the gateway "
"runs on a headless Linux box but the user's Chrome with a signed-in Google "
"profile lives on their Mac). Pass 'auto' to use the single registered node. "
"Default: run locally. Nodes are approved via `hermes meet node approve`.")},
required=["url"])
MEET_STATUS_SCHEMA: Dict[str, Any] = {
"name": "meet_status",
"description": (
"Report the current Meet session state — whether the bot is alive, has joined, is sitting "
"in the lobby, number of transcript lines captured, and last-caption timestamp."),
"parameters": {"type": "object", "properties": {"node": _NODE_PROP}, "additionalProperties": False},
}
MEET_STATUS_SCHEMA = _schema(
"meet_status",
"Report the current Meet session state — whether the bot is alive, has joined, is sitting "
"in the lobby, number of transcript lines captured, and last-caption timestamp.",
{"node": _NODE_PROP})
MEET_TRANSCRIPT_SCHEMA: Dict[str, Any] = {
"name": "meet_transcript",
"description": (
"Read the scraped transcript for the active Meet session. Returns "
"full transcript unless 'last' is set, in which case returns the last N lines only."),
"parameters": {
"type": "object",
"properties": {
"last": {
"type": "integer",
"description": ("Optional: return only the last N caption lines. Useful "
"for polling during a meeting without re-reading the whole transcript."),
"minimum": 1},
"node": _NODE_PROP},
"additionalProperties": False}}
MEET_TRANSCRIPT_SCHEMA = _schema(
"meet_transcript",
"Read the scraped transcript for the active Meet session. Returns "
"full transcript unless 'last' is set, in which case returns the last N lines only.",
{"last": {"type": "integer",
"description": ("Optional: return only the last N caption lines. Useful "
"for polling during a meeting without re-reading the whole transcript."),
"minimum": 1},
"node": _NODE_PROP})
MEET_LEAVE_SCHEMA: Dict[str, Any] = {
"name": "meet_leave",
"description": (
"Leave the active Meet call cleanly, stop caption scraping, and finalize the transcript "
"file. Safe to call when no meeting is active — returns ok=false with a reason."),
"parameters": {"type": "object", "properties": {"node": _NODE_PROP}, "additionalProperties": False},
}
MEET_LEAVE_SCHEMA = _schema(
"meet_leave",
"Leave the active Meet call cleanly, stop caption scraping, and finalize the transcript "
"file. Safe to call when no meeting is active — returns ok=false with a reason.",
{"node": _NODE_PROP})
MEET_SAY_SCHEMA: Dict[str, Any] = {
"name": "meet_say",
"description": (
"Speak text into the active Meet call. Requires the active meeting to have been joined "
"with mode='realtime'. The text is queued to the bot's OpenAI Realtime session; the "
"generated audio is streamed into Chrome's fake microphone via a virtual audio device "
"(PulseAudio null-sink on Linux, BlackHole on macOS). Returns immediately — the actual "
"speech lags by a couple of seconds."),
"parameters": {
"type": "object",
"properties": {"text": {"type": "string", "description": "Text to speak."}, "node": _NODE_PROP},
"required": ["text"],
"additionalProperties": False}}
MEET_SAY_SCHEMA = _schema(
"meet_say",
"Speak text into the active Meet call. Requires the active meeting to have been joined "
"with mode='realtime'. The text is queued to the bot's OpenAI Realtime session; the "
"generated audio is streamed into Chrome's fake microphone via a virtual audio device "
"(PulseAudio null-sink on Linux, BlackHole on macOS). Returns immediately — the actual "
"speech lags by a couple of seconds.",
{"text": _str("Text to speak."), "node": _NODE_PROP},
required=["text"])
def _json(obj: Any) -> str:

View File

@@ -45,12 +45,9 @@ def iter_plugin_dirs(root: Path) -> List[Path]:
def read_plugin_description(plugin_dir: Path) -> str:
"""Return ``description`` from ``plugin.yaml`` (empty string if absent/unreadable)."""
yaml_file = plugin_dir / "plugin.yaml"
if not yaml_file.exists():
return ""
try:
import yaml
with open(yaml_file, encoding="utf-8-sig") as f:
with open(plugin_dir / "plugin.yaml", encoding="utf-8-sig") as f:
meta = yaml.safe_load(f) or {}
return meta.get("description", "")
except Exception:
@@ -69,8 +66,8 @@ 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 stays on failure; callers needing a clean retry pop it themselves."""
"""Exec a ``_new_module`` module (None -> False); False + debug-log if it raised. The sys.modules
entry stays on failure; callers needing a clean retry pop it themselves."""
if mod is None:
return False
try:
@@ -85,11 +82,9 @@ def _exec(mod: Any, logger: Optional[logging.Logger] = None) -> bool:
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: 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.
"""
siblings onto module — the shape normal imports produce, which monkeypatch relies on."""
init_file = plugin_dir / "__init__.py"
if not init_file.exists():
return None
@@ -128,8 +123,7 @@ def load_plugin_module(module_name: str, plugin_dir: Path, *, parents: Tuple[str
class NoopPluginContext:
"""Base for 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: no-op registrations except the one a subclass overrides."""
def _noop(self, *args, **kwargs):
pass
@@ -173,8 +167,6 @@ def probe_availability(load: Callable[[], Optional[Any]]) -> bool:
"""True iff *load()* returns an instance whose ``is_available()`` (if any) is truthy."""
try:
instance = load()
if instance is None:
return False
return instance.is_available() if hasattr(instance, "is_available") else True
return instance is not None and (instance.is_available() if hasattr(instance, "is_available") else True)
except Exception:
return False