1 Commits

Author SHA1 Message Date
wouser 5be3f78fe6 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
2026-07-15 10:51:47 +02:00
+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). 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 import os
@@ -13,6 +13,7 @@ import time
import logging import logging
import json import json
import asyncio import asyncio
import subprocess
import signal import signal
from concurrent.futures import ThreadPoolExecutor from concurrent.futures import ThreadPoolExecutor
from functools import wraps from functools import wraps
@@ -28,7 +29,6 @@ from quart import Quart, request, jsonify
# ── Third-party optimizations ─────────────────────────────────────────── # ── Third-party optimizations ───────────────────────────────────────────
import httpx import httpx
import diskcache import diskcache
from tenacity import retry, stop_after_attempt, wait_exponential, retry_if_exception_type
# ── Config ────────────────────────────────────────────────────────────── # ── Config ──────────────────────────────────────────────────────────────
TIMEOUT = int(os.environ.get("HERMES_RELAY_TIMEOUT", "120")) # max seconden voor Hermes query 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 = threading.Thread(target=_prewarm_agent, daemon=True)
_prewarm_thread.start() _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: async def _get_http_client() -> httpx.AsyncClient:
"""Get or create the shared HTTP client with connection pooling.""" """Get or create the shared HTTP client with connection pooling."""
global _http_client global _http_client
@@ -150,77 +130,41 @@ def cached(ttl: int = CACHE_TTL):
return decorator 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: async def send_telegram(text: str) -> bool:
"""Stuur bericht via HTTP API met connection pooling en circuit breaker.""" """Stuur bericht via hermes send CLI (snel & robuust)."""
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
try: 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: 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 return False
@@ -228,7 +172,7 @@ async def send_telegram(text: str) -> bool:
@app.route("/ask", methods=["POST"]) @app.route("/ask", methods=["POST"])
async def ask(): async def ask():
"""Main relay endpoint — retourneert alleen het antwoord.""" """Main relay endpoint -- retourneert alleen het antwoord."""
start = time.time() start = time.time()
# Parse input # Parse input