Compare commits
1 Commits
ce9e91b1b9
..
v3.1.0
| Author | SHA1 | Date | |
|---|---|---|---|
| 5be3f78fe6 |
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user