refactor(plugins): _BotState dict init, single-phase realtime bring-up, iter_plugin_dirs reuse for user cron providers, suppress/walrus collapses
This commit is contained in:
@@ -33,10 +33,7 @@ def _is_cron_provider_dir(path: Path) -> bool:
|
||||
def _user_provider_dirs() -> List[Path]:
|
||||
"""User-installed ``$HERMES_HOME/plugins/<name>/`` dirs that look like cron providers."""
|
||||
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 [c for c in _loader.iter_plugin_dirs(user_dir) if _is_cron_provider_dir(c)] if user_dir else []
|
||||
|
||||
|
||||
def _iter_provider_dirs() -> List[Tuple[str, Path]]:
|
||||
|
||||
@@ -131,10 +131,8 @@ class ChronosCronScheduler(CronScheduler):
|
||||
if j.get("enabled") and j.get("next_run_at") and j.get("state") != "paused"}
|
||||
observed = self._list_armed()
|
||||
for job_id, fire_at in desired.items():
|
||||
if observed.get(job_id) != fire_at:
|
||||
job = get_job(job_id)
|
||||
if job:
|
||||
self._arm_logged(job, f"arm job {job_id}")
|
||||
if observed.get(job_id) != fire_at and (job := get_job(job_id)):
|
||||
self._arm_logged(job, f"arm job {job_id}")
|
||||
for job_id in observed.keys() - desired.keys():
|
||||
try:
|
||||
self._cancel(job_id)
|
||||
|
||||
@@ -4,6 +4,7 @@ Wire contract: ``docs/chronos-managed-cron-contract.md``."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import contextlib
|
||||
import logging
|
||||
from typing import Any, Dict, List
|
||||
|
||||
@@ -41,10 +42,9 @@ class NasCronClient:
|
||||
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]}")
|
||||
try:
|
||||
with contextlib.suppress(Exception):
|
||||
return resp.json() if resp.content else {}
|
||||
except Exception:
|
||||
return {}
|
||||
return {}
|
||||
|
||||
def provision(self, *, job_id: str, fire_at: str, agent_callback_url: str,
|
||||
dedup_key: str) -> Dict[str, Any]:
|
||||
|
||||
@@ -70,22 +70,15 @@ class _BotState:
|
||||
"""Single-process mutable state, flushed to ``status.json`` on each change."""
|
||||
|
||||
def __init__(self, out_dir: Path, meeting_id: str, url: str):
|
||||
for _, attr, default in _STATUS_FIELDS:
|
||||
if attr:
|
||||
setattr(self, attr, default)
|
||||
self.out_dir = out_dir
|
||||
self.meeting_id = meeting_id
|
||||
self.url = url
|
||||
self._seen: set = set() # "speaker|text" keys already written
|
||||
self.__dict__.update({attr: default for _, attr, default in _STATUS_FIELDS if attr})
|
||||
self.__dict__.update(out_dir=out_dir, meeting_id=meeting_id, url=url, _seen=set(), # seen "speaker|text"
|
||||
transcript_path=out_dir / "transcript.txt", status_path=out_dir / "status.json")
|
||||
out_dir.mkdir(parents=True, exist_ok=True)
|
||||
self.transcript_path = out_dir / "transcript.txt"
|
||||
self.status_path = out_dir / "status.json"
|
||||
self._flush()
|
||||
|
||||
def record_caption(self, speaker: str, text: str) -> None:
|
||||
"""Append a caption line unless this exact (speaker, text) was already seen."""
|
||||
speaker = (speaker or "").strip() or "Unknown"
|
||||
text = (text or "").strip()
|
||||
speaker, text = (speaker or "").strip() or "Unknown", (text or "").strip()
|
||||
key = f"{speaker}|{text}"
|
||||
if not text or key in self._seen:
|
||||
return
|
||||
@@ -238,21 +231,19 @@ def _start_pcm_pump(rt: dict, bridge_info: dict, pcm_path: Path, state: "_BotSta
|
||||
|
||||
def _start_realtime_speaker(rt: dict, cfg: "_BotConfig", stop_flag: dict, state: "_BotState") -> None:
|
||||
"""Wire up the OpenAI Realtime session, the say-queue speaker thread and the PCM pump."""
|
||||
try:
|
||||
from plugins.google_meet.realtime.openai_client import RealtimeSession, RealtimeSpeaker
|
||||
except Exception as e:
|
||||
state.set(error=f"realtime import failed: {e}")
|
||||
return
|
||||
pcm_path, queue_path = cfg.out_dir / "speaker.pcm", 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
|
||||
phase = "import"
|
||||
try:
|
||||
from plugins.google_meet.realtime.openai_client import RealtimeSession, RealtimeSpeaker
|
||||
phase = "connect"
|
||||
session = RealtimeSession(
|
||||
api_key=cfg.realtime_api_key, model=cfg.realtime_model, voice=cfg.realtime_voice,
|
||||
instructions=cfg.realtime_instructions, audio_sink_path=pcm_path, sample_rate=24000)
|
||||
session.connect()
|
||||
except Exception as e:
|
||||
state.set(error=f"realtime connect failed: {e}")
|
||||
state.set(error=f"realtime {phase} failed: {e}")
|
||||
return
|
||||
rt["session"] = session
|
||||
speaker = RealtimeSpeaker(session=session, queue_path=queue_path,
|
||||
|
||||
Reference in New Issue
Block a user