From 5be3f78fe60faff265ebde924a68f27dd87d80e4 Mon Sep 17 00:00:00 2001 From: wouser Date: Wed, 15 Jul 2026 10:51:47 +0200 Subject: [PATCH] v3.1: Fix Telegram delivery via hermes send CLI subprocess - Replace send_message_tool import with hermes send CLI subprocess - Remove broken _ensure_gw_config() that clobbered env vars - Remove unused CircuitBreaker class and tenacity import - Clean up docstring and unused imports - Token fix: requires real TELEGRAM_BOT_TOKEN in ~/.hermes/.env - Systemd drop-in: /etc/systemd/system/hermes-relay.service.d/telegram.conf --- app.py | 130 ++++++++++++++++----------------------------------------- 1 file changed, 37 insertions(+), 93 deletions(-) diff --git a/app.py b/app.py index 487c440..b3021eb 100644 --- a/app.py +++ b/app.py @@ -1,10 +1,10 @@ """ -Hermes Relay v3 — Quart (async) bridge tussen Node-RED en Hermes Agent. +Hermes Relay v3.1 - Quart (async) bridge tussen Node-RED en Hermes Agent. Gebruikt de voice-assistant profiel (minimale chatbot, max_turns=10). -Geen prefix/model override nodig — profiel handelt dit af. +Geen prefix/model override nodig, profiel handelt dit af. -Node-RED stuurt POST met msg.payload → Hermes (direct) → Alleen antwoord terug. +Node-RED stuurt POST met msg.payload -> Hermes (direct) -> Alleen antwoord terug. """ import os @@ -13,6 +13,7 @@ import time import logging import json import asyncio +import subprocess import signal from concurrent.futures import ThreadPoolExecutor from functools import wraps @@ -28,7 +29,6 @@ from quart import Quart, request, jsonify # ── Third-party optimizations ─────────────────────────────────────────── import httpx import diskcache -from tenacity import retry, stop_after_attempt, wait_exponential, retry_if_exception_type # ── Config ────────────────────────────────────────────────────────────── TIMEOUT = int(os.environ.get("HERMES_RELAY_TIMEOUT", "120")) # max seconden voor Hermes query @@ -84,26 +84,6 @@ def _prewarm_agent(): _prewarm_thread = threading.Thread(target=_prewarm_agent, daemon=True) _prewarm_thread.start() -# Laad gateway config eenmalig voor Telegram verzenden -_gw_config_loaded = False -_gw_config = None - -def _ensure_gw_config(): - """Laad gateway config eenmalig.""" - global _gw_config_loaded, _gw_config - if _gw_config_loaded: - return _gw_config - # Forceer default profiel voor gateway config - os.environ["HERMES_HOME"] = "/root/.hermes" - # Telegram token en channel uit environment variables (niet hardcoded) - os.environ["TELEGRAM_BOT_TOKEN"] = os.environ.get("TELEGRAM_BOT_TOKEN", "") - os.environ["TELEGRAM_HOME_CHANNEL"] = os.environ.get("TELEGRAM_HOME_CHANNEL", "") - from gateway.config import load_gateway_config - _gw_config = load_gateway_config() - _gw_config_loaded = True - return _gw_config - - async def _get_http_client() -> httpx.AsyncClient: """Get or create the shared HTTP client with connection pooling.""" global _http_client @@ -150,77 +130,41 @@ def cached(ttl: int = CACHE_TTL): return decorator -class CircuitBreaker: - """Simple circuit breaker for external service calls.""" - def __init__(self, failure_threshold: int = 5, recovery_timeout: int = 60): - self.failure_threshold = failure_threshold - self.recovery_timeout = recovery_timeout - self.failure_count = 0 - self.last_failure_time: Optional[float] = None - self.state = "closed" # closed, open, half-open - self._lock = asyncio.Lock() - - async def call(self, func: Callable, *args, **kwargs) -> Any: - async with self._lock: - if self.state == "open": - if self.last_failure_time is not None and time.time() - self.last_failure_time > self.recovery_timeout: - self.state = "half-open" - logger.info("Circuit breaker moving to half-open state") - else: - raise Exception("Circuit breaker is open") - - try: - result = await func(*args, **kwargs) - async with self._lock: - if self.state == "half-open": - self.state = "closed" - self.failure_count = 0 - logger.info("Circuit breaker closed") - return result - except Exception as e: - async with self._lock: - self.failure_count += 1 - self.last_failure_time = time.time() - if self.failure_count >= self.failure_threshold: - self.state = "open" - logger.warning("Circuit breaker opened after %d failures", self.failure_count) - raise - - -# Circuit breaker for Telegram -_tg_circuit_breaker = CircuitBreaker(failure_threshold=5, recovery_timeout=60) - - async def send_telegram(text: str) -> bool: - """Stuur bericht via HTTP API met connection pooling en circuit breaker.""" - async def _do_send(): - try: - logger.info("Telegram send: starting...") - _ensure_gw_config() - logger.info("Telegram send: gw config loaded") - from tools.send_message_tool import send_message_tool - - # Stuur via home channel (telegram = default target) - result_json = send_message_tool({ - "action": "send", - "target": TG_TARGET, - "message": text, - }) - logger.info("Telegram send: send_message_tool returned") - result = json.loads(result_json) - if result.get("error"): - logger.error("Telegram send failed: %s", result["error"]) - return False - logger.info("Telegram send: success=%s", result.get("success", False)) - return result.get("success", False) - except Exception as e: - logger.error("Telegram send failed with exception: %s", e, exc_info=True) - return False - + """Stuur bericht via hermes send CLI (snel & robuust).""" try: - return await _tg_circuit_breaker.call(_do_send) + logger.info("Telegram send: starting via hermes send CLI...") + result = subprocess.run( + [ + "/usr/local/lib/hermes-agent/venv/bin/hermes", "send", + "--to", TG_TARGET, + "--quiet", + text, + ], + capture_output=True, + text=True, + timeout=15, + env={**os.environ, "HERMES_HOME": "/root/.hermes"}, + ) + if result.returncode == 0: + logger.info("Telegram send: success") + return True + else: + stderr = result.stderr.strip() + # "Skipped" message (cron detection false positive) is not a failure + if "skipped" in stderr.lower() or "cron" in stderr.lower(): + logger.info("Telegram send: skipped by hermes send (non-fatal)") + return True + logger.error("Telegram send failed: %s", stderr[:200]) + return False + except subprocess.TimeoutExpired: + logger.error("Telegram send: timeout") + return False + except FileNotFoundError: + logger.error("Telegram send: hermes binary not found") + return False except Exception as e: - logger.error("Telegram send blocked by circuit breaker: %s", e) + logger.error("Telegram send failed with exception: %s", e) return False @@ -228,7 +172,7 @@ async def send_telegram(text: str) -> bool: @app.route("/ask", methods=["POST"]) async def ask(): - """Main relay endpoint — retourneert alleen het antwoord.""" + """Main relay endpoint -- retourneert alleen het antwoord.""" start = time.time() # Parse input