""" 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")