5be3f78fe6
- 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
280 lines
9.9 KiB
Python
280 lines
9.9 KiB
Python
"""
|
|
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.
|
|
|
|
Node-RED stuurt POST met msg.payload -> Hermes (direct) -> Alleen antwoord terug.
|
|
"""
|
|
|
|
import os
|
|
import sys
|
|
import time
|
|
import logging
|
|
import json
|
|
import asyncio
|
|
import subprocess
|
|
import signal
|
|
from concurrent.futures import ThreadPoolExecutor
|
|
from functools import wraps
|
|
from typing import Optional, Callable, Any
|
|
|
|
# ── Hermes path toevoegen (zodat we Hermes direct kunnen importeren) ───
|
|
_HERMES_SITE = os.environ.get("HERMES_RELAY_SITE_PACKAGES", "/usr/local/lib/hermes-agent/venv/lib/python3.11/site-packages")
|
|
if _HERMES_SITE not in sys.path:
|
|
sys.path.insert(0, _HERMES_SITE)
|
|
|
|
from quart import Quart, request, jsonify
|
|
|
|
# ── Third-party optimizations ───────────────────────────────────────────
|
|
import httpx
|
|
import diskcache
|
|
|
|
# ── Config ──────────────────────────────────────────────────────────────
|
|
TIMEOUT = int(os.environ.get("HERMES_RELAY_TIMEOUT", "120")) # max seconden voor Hermes query
|
|
TG_TARGET = os.environ.get("HERMES_RELAY_TELEGRAM_TARGET", "telegram") # hermes send target
|
|
|
|
# Thread pool config - increased for better concurrency
|
|
THREAD_POOL_WORKERS = int(os.environ.get("HERMES_RELAY_THREAD_WORKERS", "8"))
|
|
|
|
# Cache config
|
|
CACHE_DIR = os.environ.get("HERMES_RELAY_CACHE_DIR", "/tmp/hermes-relay-cache")
|
|
CACHE_TTL = int(os.environ.get("HERMES_RELAY_CACHE_TTL", "300")) # 5 minuten default
|
|
|
|
# HTTP client config for Telegram
|
|
HTTP_CLIENT_TIMEOUT = httpx.Timeout(10.0, connect=5.0)
|
|
HTTP_CLIENT_LIMITS = httpx.Limits(max_keepalive_connections=10, max_connections=20)
|
|
|
|
# ── Logging ─────────────────────────────────────────────────────────────
|
|
logging.basicConfig(
|
|
level=logging.INFO,
|
|
format="[hermes-relay] %(asctime)s %(levelname)s %(message)s",
|
|
stream=sys.stderr,
|
|
)
|
|
logger = logging.getLogger(__name__)
|
|
|
|
# ── Quart app ───────────────────────────────────────────────────────────
|
|
app = Quart(__name__)
|
|
|
|
# Thread pool for running sync Hermes agent calls
|
|
_executor = ThreadPoolExecutor(max_workers=THREAD_POOL_WORKERS)
|
|
|
|
# HTTP client with connection pooling for Telegram
|
|
_http_client: Optional[httpx.AsyncClient] = None
|
|
|
|
# Disk cache for health checks and other reusable data
|
|
_cache = diskcache.Cache(CACHE_DIR, timeout=5)
|
|
|
|
# ── Hermes warm houden ──────────────────────────────────────────────────
|
|
from hermes_client import call_hermes as _hermes_query
|
|
from hermes_client import _ensure_hermes_loaded, _initialize_agent
|
|
|
|
# Pre-warm the agent on startup (non-blocking)
|
|
import threading
|
|
|
|
def _prewarm_agent():
|
|
"""Pre-initialize Hermes agent in background so first request is fast."""
|
|
try:
|
|
_ensure_hermes_loaded()
|
|
_initialize_agent()
|
|
logger.info("Hermes agent pre-warmed successfully")
|
|
except Exception as e:
|
|
logger.error("Hermes agent pre-warm failed: %s", e, exc_info=True)
|
|
|
|
_prewarm_thread = threading.Thread(target=_prewarm_agent, daemon=True)
|
|
_prewarm_thread.start()
|
|
|
|
async def _get_http_client() -> httpx.AsyncClient:
|
|
"""Get or create the shared HTTP client with connection pooling."""
|
|
global _http_client
|
|
if _http_client is None or _http_client.is_closed:
|
|
_http_client = httpx.AsyncClient(
|
|
timeout=HTTP_CLIENT_TIMEOUT,
|
|
limits=HTTP_CLIENT_LIMITS,
|
|
http2=True, # Enable HTTP/2 for better performance
|
|
)
|
|
return _http_client
|
|
|
|
|
|
async def _close_http_client():
|
|
"""Close the HTTP client on shutdown."""
|
|
global _http_client
|
|
if _http_client and not _http_client.is_closed:
|
|
await _http_client.aclose()
|
|
_http_client = None
|
|
|
|
|
|
def cached(ttl: int = CACHE_TTL):
|
|
"""Decorator for caching function results in diskcache."""
|
|
def decorator(func: Callable) -> Callable:
|
|
@wraps(func)
|
|
async def wrapper(*args, **kwargs):
|
|
# Create cache key from function name and arguments
|
|
key_parts = [func.__name__]
|
|
key_parts.extend(str(arg) for arg in args)
|
|
key_parts.extend(f"{k}={v}" for k, v in sorted(kwargs.items()))
|
|
cache_key = ":".join(key_parts)
|
|
|
|
# Try to get from cache
|
|
cached_result = _cache.get(cache_key)
|
|
if cached_result is not None:
|
|
logger.debug("Cache hit for %s", cache_key)
|
|
return cached_result
|
|
|
|
# Call function and cache result
|
|
result = await func(*args, **kwargs)
|
|
_cache.set(cache_key, result, expire=ttl)
|
|
logger.debug("Cache set for %s", cache_key)
|
|
return result
|
|
return wrapper
|
|
return decorator
|
|
|
|
|
|
async def send_telegram(text: str) -> bool:
|
|
"""Stuur bericht via hermes send CLI (snel & robuust)."""
|
|
try:
|
|
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 failed with exception: %s", e)
|
|
return False
|
|
|
|
|
|
# ── Routes ──────────────────────────────────────────────────────────────
|
|
|
|
@app.route("/ask", methods=["POST"])
|
|
async def ask():
|
|
"""Main relay endpoint -- retourneert alleen het antwoord."""
|
|
start = time.time()
|
|
|
|
# Parse input
|
|
if request.is_json:
|
|
data = await request.get_json(silent=True) or {}
|
|
prompt = data.get("payload", "")
|
|
else:
|
|
form = await request.form
|
|
prompt = form.get("payload", "")
|
|
|
|
if not prompt:
|
|
return jsonify({"status": "error", "message": "No payload received"}), 400
|
|
|
|
logger.info("Received: %s...", prompt[:100])
|
|
|
|
# Roep Hermes aan (in thread pool omdat Hermes sync is)
|
|
loop = asyncio.get_event_loop()
|
|
answer = await loop.run_in_executor(_executor, _hermes_query, prompt)
|
|
elapsed = round(time.time() - start, 1)
|
|
|
|
logger.info("Answer (%.1fs): %s...", elapsed, answer[:200])
|
|
|
|
# Stuur ook naar Telegram (alleen het antwoord) - in thread pool
|
|
tg_sent = await loop.run_in_executor(_executor, lambda: asyncio.run(send_telegram(answer)))
|
|
|
|
return jsonify({
|
|
"status": "ok",
|
|
"answer": answer,
|
|
"elapsed_seconds": elapsed,
|
|
"telegram_sent": tg_sent,
|
|
})
|
|
|
|
|
|
@app.route("/health", methods=["GET"])
|
|
async def health():
|
|
"""Health check voor Node-RED monitoring.
|
|
|
|
Query params:
|
|
- detail=true: run deep health check (includes LLM query test)
|
|
"""
|
|
from hermes_client import health_check
|
|
|
|
detail = request.args.get("detail", "false").lower() == "true"
|
|
|
|
if detail:
|
|
# Use cached health check for detail view
|
|
checks = await _cached_health_check()
|
|
return jsonify({
|
|
"status": "ok" if checks["healthy"] else "degraded",
|
|
"service": "hermes-relay",
|
|
"version": "3.0",
|
|
"mode": "async-quart",
|
|
"profile": "voice-assistant",
|
|
"checks": checks,
|
|
})
|
|
|
|
# Quick health check - just return basic status (cached)
|
|
quick_check = await _cached_quick_health()
|
|
return jsonify(quick_check)
|
|
|
|
|
|
@cached(ttl=60) # Cache quick health for 60 seconds
|
|
async def _cached_quick_health():
|
|
"""Cached quick health check."""
|
|
return {
|
|
"status": "ok",
|
|
"service": "hermes-relay",
|
|
"version": "3.0",
|
|
"mode": "async-quart",
|
|
"profile": "voice-assistant",
|
|
}
|
|
|
|
|
|
@cached(ttl=300) # Cache detailed health for 5 minutes
|
|
async def _cached_health_check():
|
|
"""Cached detailed health check."""
|
|
from hermes_client import health_check
|
|
checks = health_check()
|
|
return checks
|
|
|
|
|
|
@app.before_serving
|
|
async def startup():
|
|
"""Startup tasks."""
|
|
# Ensure thread pool is ready
|
|
# Pre-warm HTTP client
|
|
await _get_http_client()
|
|
logger.info("Hermes Relay v3.0 startup complete (workers=%d, http2=enabled)", THREAD_POOL_WORKERS)
|
|
|
|
|
|
@app.after_serving
|
|
async def shutdown():
|
|
"""Cleanup on shutdown."""
|
|
_executor.shutdown(wait=True)
|
|
await _close_http_client()
|
|
_cache.close()
|
|
logger.info("Hermes Relay v3.0 shutdown complete")
|
|
|
|
|
|
if __name__ == "__main__":
|
|
import asyncio
|
|
logger.info("Starting Hermes Relay v3.0 (voice-assistant profiel, Quart async)...")
|
|
# Run with uvicorn for production
|
|
import uvicorn
|
|
uvicorn.run(app, host="0.0.0.0", port=8650, log_level="info") |