From 67160b2de0db52aa4762e16efb83071ed762812c Mon Sep 17 00:00:00 2001 From: wouser Date: Mon, 20 Jul 2026 06:17:31 +0200 Subject: [PATCH] fix: make relay fast and concurrency-safe --- CHANGELOG.md | 13 ++ README.md | 118 ++++++++++-------- app.py | 316 +++++++++++++++++++---------------------------- hermes_client.py | 8 +- requirements.txt | 18 +-- 5 files changed, 221 insertions(+), 252 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index e2d4b3b..849470e 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -5,6 +5,19 @@ All notable changes to Hermes Relay will be documented in this file. The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/), and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html). +## [3.2.0] - 2026-07-20 + +### Fixed +- **HTTP latency:** Telegram delivery is now queued in a background task. `/ask` returns after the LLM response instead of waiting another ~5–6 seconds for `hermes send`. +- **Shared-agent concurrency:** model calls are serialized with an async lock. The warm Hermes agent mutates message state and is not safe for concurrent calls. +- **Timeout enforcement:** `HERMES_RELAY_TIMEOUT` now covers queueing plus the model call and returns HTTP `504` on expiry. +- **Accurate API metrics:** `elapsed_seconds` / `llm_elapsed_seconds` report model duration, and `request_elapsed_seconds` reports actual HTTP duration. `telegram_queued` replaces the misleading synchronous `telegram_sent` result. +- **Toolset warning:** default relay toolsets are empty; Telegram delivery uses `hermes send` and does not require a `messaging` toolset. +- **Health semantics/version:** unified version reporting at 3.2.0 and clarified that detailed health checks local agent state rather than performing an LLM probe. + +### Removed +- Unused `httpx` HTTP/2 client, disk cache and related dependencies. Telegram delivery is performed by the Hermes CLI subprocess, not this client. + ## [3.1.0] - 2026-07-15 ### Fixed diff --git a/README.md b/README.md index f8d34f2..e62722e 100644 --- a/README.md +++ b/README.md @@ -1,75 +1,95 @@ -# Hermes Relay +# Hermes Relay v3.2 -Flask bridge tussen Node-RED en Hermes Agent. - -## Architectuur +Snelle Quart-bridge tussen Node-RED en een warme Hermes Agent. +```text +Node-RED ──POST /ask──▶ Quart relay ──serialized call──▶ warme Hermes-agent + │ + ├── JSON direct terug naar Node-RED + └── Telegram-bezorging op de achtergrond via `hermes send` ``` -Node-RED (192.168.1.125) ──POST /ask──▶ Hermes Relay (192.168.1.74:8650) - │ - ├──▶ hermes chat -q (subprocess) - │ │ - │ ▼ - │ JSON response ──▶ Node-RED - │ - └──▶ hermes send --to telegram - │ - ▼ - Telegram chat -``` + +## Gedrag en grenzen + +- **Warme agent:** Hermes wordt tijdens startup geladen/geïnitialiseerd. +- **Veilige concurrency:** de gedeelde, mutable agent verwerkt precies één LLM-call tegelijk. Overige HTTP-verzoeken wachten in de queue; ze kunnen niet elkaars `messages` resetten. +- **Timeout:** `HERMES_RELAY_TIMEOUT` (standaard 120 s) omvat wachttijd plus LLM-call. Een timeout retourneert HTTP `504`. +- **Snelle response:** Telegram ligt niet meer in het HTTP-pad. Een succesvolle `/ask` reageert zodra het model antwoordt. +- **Background delivery:** Telegram wordt betrouwbaar gelogd. Bij gecontroleerde shutdown wacht de service op lopende bezorgingen. ## Endpoints ### `POST /ask` -Verwacht JSON body: + +JSON body: + ```json { "payload": "Wat is het weer in Best?" } ``` -Response: +Succesresponse: + ```json { "status": "ok", - "answer": "Het is 22°C en zonnig in Best...", - "elapsed_seconds": 8.3, - "telegram_sent": true + "answer": "...", + "elapsed_seconds": 6.213, + "llm_elapsed_seconds": 6.213, + "request_elapsed_seconds": 6.214, + "telegram_queued": true } ``` -### `GET /health` -Health check voor monitoring. +`elapsed_seconds` blijft aanwezig voor compatibiliteit en is de LLM-duur. `request_elapsed_seconds` is de feitelijke HTTP-duur. `telegram_queued` betekent dat bezorging geaccepteerd is voor de achtergrondtaak; raadpleeg journald voor het uiteindelijke bezorgresultaat. -## Setup +Fouten: `400` lege payload, `413` payload groter dan limiet, `504` timeout, `502` onverwachte agentfout. + +### `GET /health` + +Geeft service-status, versie en het aantal lopende Telegram-bezorgingen terug. + +- `/health?detail=true` voegt de **lokale agent-state** toe; dit doet bewust geen dure LLM/provider-call. + +## Configuratie + +| Variabele | Standaard | Betekenis | +|---|---:|---| +| `HERMES_RELAY_TIMEOUT` | `120` | Max. wachttijd + modelcall per HTTP-request | +| `HERMES_RELAY_DELIVERY_WORKERS` | `2` | Begrensde workers voor Telegram-bezorging | +| `HERMES_RELAY_MAX_PAYLOAD_CHARS` | `12000` | Maximale lengte van `payload` | +| `HERMES_RELAY_MAX_TURNS` | `10` | Agent turn-budget | +| `HERMES_RELAY_MODEL` | profieldefault | Optionele modelovertuiging | +| `HERMES_RELAY_TOOLSETS` | leeg | Alleen invullen met geldige benodigde toolsets | + +De relay leest zijn Hermes-config via profiel `voice-assistant`. Pas model/provider/config aan en herstart daarna de service: ```bash -cd /root/hermes-relay -python3 -m venv venv -source venv/bin/activate -pip install -r requirements.txt - -# Test -python app.py - -# Of als service -sudo systemctl enable hermes-relay -sudo systemctl start hermes-relay +systemctl restart hermes-relay ``` -## Node-RED Nodes +## Testen -### Versturen (HTTP Request node) -- **Method:** POST +```bash +curl -s http://127.0.0.1:8650/health +curl -s -X POST http://127.0.0.1:8650/ask \ + -H 'Content-Type: application/json' \ + -d '{"payload":"Zeg alleen hoi"}' + +journalctl -u hermes-relay -f +``` + +## Node-RED + +Gebruik een HTTP Request-node: + +- **Method:** `POST` - **URL:** `http://192.168.1.74:8650/ask` -- **Return:** a parsed JSON object -- **Timeout:** 120000 (ms) — Hermes kan even duren +- **Return:** parsed JSON +- **Timeout:** `120000` ms -### Ontvangen -Het antwoord komt terug als `msg.payload`: -- `msg.payload.answer` — het Hermes antwoord -- `msg.payload.elapsed_seconds` — hoe lang het duurde -- `msg.payload.telegram_sent` — of het ook naar Telegram is gestuurd -- `msg.payload.status` — "ok" of "error" +Voorafgaande Function-node: -### Health check (optioneel) -- **Method:** GET -- **URL:** `http://192.168.1.74:8650/health` +```javascript +msg.payload = { payload: msg.payload }; +return msg; +``` diff --git a/app.py b/app.py index b3021eb..9550132 100644 --- a/app.py +++ b/app.py @@ -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") \ No newline at end of file + 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") diff --git a/hermes_client.py b/hermes_client.py index d3a9bb8..6faec25 100644 --- a/hermes_client.py +++ b/hermes_client.py @@ -23,7 +23,13 @@ logger = logging.getLogger(__name__) HERMES_HOME = os.environ.get("HERMES_RELAY_HERMES_HOME", "/root/.hermes/profiles/voice-assistant") HERMES_SITE_PACKAGES = os.environ.get("HERMES_RELAY_SITE_PACKAGES", "/usr/local/lib/hermes-agent/venv/lib/python3.11/site-packages") HERMES_MODEL = os.environ.get("HERMES_RELAY_MODEL", "") # Empty = use profile default -HERMES_TOOLSETS = os.environ.get("HERMES_RELAY_TOOLSETS", "messaging").split(",") +# Telegram delivery runs via `hermes send`, so the relay agent needs no +# messaging toolset. Filter empty items to avoid an unknown-toolset warning. +HERMES_TOOLSETS = [ + tool.strip() + for tool in os.environ.get("HERMES_RELAY_TOOLSETS", "").split(",") + if tool.strip() +] HERMES_MAX_TURNS = int(os.environ.get("HERMES_RELAY_MAX_TURNS", "10")) HERMES_TIMEOUT = int(os.environ.get("HERMES_RELAY_TIMEOUT", "120")) LOG_LEVEL = os.environ.get("HERMES_RELAY_LOG_LEVEL", "INFO").upper() diff --git a/requirements.txt b/requirements.txt index a77eedb..5229fca 100644 --- a/requirements.txt +++ b/requirements.txt @@ -1,19 +1,9 @@ -# Core dependencies - pinned versions for reproducible builds -flask==3.0.3 -requests==2.32.3 - -# Async/Production dependencies +# Web service quart==0.19.4 uvicorn[standard]==0.30.6 -gunicorn==22.0.0 -# Performance optimizations -httpx==0.27.2 # HTTP/2 connection pooling for Telegram -diskcache==5.6.3 # Semantic caching -brotli==1.1.0 # Compression - -# Resilience patterns -tenacity==8.5.0 # Retry logic with exponential backoff +# Optional production server +hypercorn==0.17.3 # Type hints support -typing_extensions==4.12.2 \ No newline at end of file +typing_extensions==4.12.2