fix: make relay fast and concurrency-safe

This commit is contained in:
2026-07-20 06:17:31 +02:00
parent 4eaa21f20c
commit 67160b2de0
5 changed files with 221 additions and 252 deletions
+128 -188
View File
@@ -1,51 +1,30 @@
"""
Hermes Relay v3.1 - Quart (async) bridge tussen Node-RED en Hermes Agent.
"""Hermes Relay v3.2 — fast, bounded Quart bridge for Node-RED."""
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 logging
import os
import subprocess
import signal
import sys
import threading
import time
from concurrent.futures import ThreadPoolExecutor
from functools import wraps
from typing import Optional, Callable, Any
from typing import Set
# ── 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")
_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
from quart import Quart, jsonify, request
# ── Third-party optimizations ───────────────────────────────────────────
import httpx
import diskcache
VERSION = "3.2.0"
TIMEOUT = int(os.environ.get("HERMES_RELAY_TIMEOUT", "120"))
TG_TARGET = os.environ.get("HERMES_RELAY_TELEGRAM_TARGET", "telegram")
DELIVERY_WORKERS = int(os.environ.get("HERMES_RELAY_DELIVERY_WORKERS", "2"))
MAX_PAYLOAD_CHARS = int(os.environ.get("HERMES_RELAY_MAX_PAYLOAD_CHARS", "12000"))
# ── 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",
@@ -53,129 +32,91 @@ logging.basicConfig(
)
logger = logging.getLogger(__name__)
# ── Quart app ───────────────────────────────────────────────────────────
app = Quart(__name__)
# Thread pool for running sync Hermes agent calls
_executor = ThreadPoolExecutor(max_workers=THREAD_POOL_WORKERS)
# HermesCLI/AIAgent is mutable (messages are reset per request), so it must not
# be called concurrently. Delivery uses a separate, bounded executor instead.
_hermes_lock: asyncio.Lock | None = None
_delivery_executor = ThreadPoolExecutor(
max_workers=DELIVERY_WORKERS, thread_name_prefix="telegram-delivery"
)
_delivery_tasks: Set[asyncio.Task] = set()
# 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."""
def _prewarm_agent() -> None:
"""Initialize the agent before the first request without blocking startup."""
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)
except Exception:
logger.exception("Hermes agent pre-warm failed")
_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)."""
def _send_telegram_sync(text: str) -> bool:
"""Send through Hermes CLI. This deliberately runs outside the request path."""
try:
logger.info("Telegram send: starting via hermes send CLI...")
logger.info("Telegram send: starting via hermes send CLI")
result = subprocess.run(
[
"/usr/local/lib/hermes-agent/venv/bin/hermes", "send",
"--to", TG_TARGET,
"/usr/local/lib/hermes-agent/venv/bin/hermes",
"send",
"--to",
TG_TARGET,
"--quiet",
text,
],
capture_output=True,
text=True,
timeout=15,
timeout=30,
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
logger.error("Telegram send failed (rc=%d): %s", result.returncode, result.stderr.strip()[:300])
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
except Exception:
logger.exception("Telegram send failed")
return False
# ── Routes ──────────────────────────────────────────────────────────────
async def _deliver_in_background(text: str) -> None:
loop = asyncio.get_running_loop()
await loop.run_in_executor(_delivery_executor, _send_telegram_sync, text)
def _track_delivery(task: asyncio.Task) -> None:
_delivery_tasks.discard(task)
try:
task.result()
except asyncio.CancelledError:
logger.warning("Telegram delivery cancelled during shutdown")
except Exception:
logger.exception("Unexpected Telegram delivery task failure")
def _queue_telegram(text: str) -> None:
task = asyncio.create_task(_deliver_in_background(text), name="hermes-relay-telegram")
_delivery_tasks.add(task)
task.add_done_callback(_track_delivery)
@app.route("/ask", methods=["POST"])
async def ask():
"""Main relay endpoint -- retourneert alleen het antwoord."""
start = time.time()
"""Answer a request; Telegram delivery is asynchronous and non-blocking."""
request_started = time.perf_counter()
# Parse input
if request.is_json:
data = await request.get_json(silent=True) or {}
prompt = data.get("payload", "")
@@ -183,98 +124,97 @@ async def ask():
form = await request.form
prompt = form.get("payload", "")
if not prompt:
if not isinstance(prompt, str) or not prompt.strip():
return jsonify({"status": "error", "message": "No payload received"}), 400
if len(prompt) > MAX_PAYLOAD_CHARS:
return jsonify({"status": "error", "message": f"Payload exceeds {MAX_PAYLOAD_CHARS} characters"}), 413
logger.info("Received: %s...", prompt[:100])
lock = _hermes_lock
if lock is None: # Defensive fallback for an early request during startup.
return jsonify({"status": "error", "message": "Service is starting; retry shortly"}), 503
# 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)
try:
# One warm agent, therefore one model call at a time. The timeout covers
# queue wait plus the call, preventing an indefinitely occupied request.
async with asyncio.timeout(TIMEOUT):
async with lock:
llm_started = time.perf_counter()
answer = await asyncio.to_thread(_hermes_query, prompt)
llm_elapsed = round(time.perf_counter() - llm_started, 3)
except TimeoutError:
request_elapsed = round(time.perf_counter() - request_started, 3)
logger.warning("Hermes request timed out after %.3fs", request_elapsed)
return jsonify({
"status": "error",
"message": "Hermes request timed out",
"request_elapsed_seconds": request_elapsed,
}), 504
except Exception:
logger.exception("Unhandled Hermes request failure")
return jsonify({"status": "error", "message": "Hermes request failed"}), 502
logger.info("Answer (%.1fs): %s...", elapsed, answer[:200])
if not isinstance(answer, str):
logger.error("Hermes returned non-string answer: %s", type(answer).__name__)
return jsonify({"status": "error", "message": "Hermes returned an invalid response"}), 502
# Stuur ook naar Telegram (alleen het antwoord) - in thread pool
tg_sent = await loop.run_in_executor(_executor, lambda: asyncio.run(send_telegram(answer)))
request_elapsed = round(time.perf_counter() - request_started, 3)
logger.info("Answer (llm=%.3fs, request=%.3fs): %s...", llm_elapsed, request_elapsed, answer[:200])
_queue_telegram(answer)
return jsonify({
"status": "ok",
"answer": answer,
"elapsed_seconds": elapsed,
"telegram_sent": tg_sent,
# Backwards-compatible field: model execution time, not delivery time.
"elapsed_seconds": llm_elapsed,
"llm_elapsed_seconds": llm_elapsed,
"request_elapsed_seconds": request_elapsed,
"telegram_queued": True,
})
@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",
"""Return relay state. `detail=true` checks local agent state, not the LLM."""
base = {
"service": "hermes-relay",
"version": "3.0",
"version": VERSION,
"mode": "async-quart",
"profile": "voice-assistant",
"pending_telegram_deliveries": len(_delivery_tasks),
}
@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
if request.args.get("detail", "false").lower() == "true":
from hermes_client import health_check
checks = await asyncio.to_thread(health_check)
return jsonify({
**base,
"status": "ok" if checks["healthy"] else "degraded",
"checks": checks,
})
return jsonify({**base, "status": "ok"})
@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)
async def startup() -> None:
global _hermes_lock
_hermes_lock = asyncio.Lock()
logger.info(
"Hermes Relay v%s startup complete (one serialized agent, delivery_workers=%d, timeout=%ds)",
VERSION, DELIVERY_WORKERS, TIMEOUT,
)
@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")
async def shutdown() -> None:
pending = tuple(_delivery_tasks)
if pending:
logger.info("Waiting for %d Telegram delivery task(s)", len(pending))
await asyncio.gather(*pending, return_exceptions=True)
_delivery_executor.shutdown(wait=True, cancel_futures=False)
logger.info("Hermes Relay v%s shutdown complete", VERSION)
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")
logger.info("Starting Hermes Relay v%s (voice-assistant profile, Quart async)", VERSION)
uvicorn.run(app, host="0.0.0.0", port=8650, log_level="info")