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
This commit is contained in:
2026-07-15 10:51:47 +02:00
parent ce9e91b1b9
commit 5be3f78fe6
+37 -93
View File
@@ -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