Files
hermes-relay/app.py
T

228 lines
8.2 KiB
Python

"""Hermes Relay v3.2 — fast, bounded Quart bridge for Node-RED."""
import asyncio
import logging
import os
import subprocess
import sys
import threading
import time
from concurrent.futures import ThreadPoolExecutor
from typing import Set
_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, jsonify, request
VERSION = "3.3.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"))
# HTTP replies are the default. Set this explicitly to true only when relay
# answers should additionally be delivered through Telegram.
TELEGRAM_DELIVERY_ENABLED = os.environ.get(
"HERMES_RELAY_TELEGRAM_DELIVERY_ENABLED", "false"
).strip().lower() in {"1", "true", "yes", "on"}
logging.basicConfig(
level=logging.INFO,
format="[hermes-relay] %(asctime)s %(levelname)s %(message)s",
stream=sys.stderr,
)
logger = logging.getLogger(__name__)
app = Quart(__name__)
# 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()
from hermes_client import call_hermes as _hermes_query
from hermes_client import _ensure_hermes_loaded, _initialize_agent
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:
logger.exception("Hermes agent pre-warm failed")
_prewarm_thread = threading.Thread(target=_prewarm_agent, daemon=True)
_prewarm_thread.start()
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")
result = subprocess.run(
[
"/usr/local/lib/hermes-agent/venv/bin/hermes",
"send",
"--to",
TG_TARGET,
"--quiet",
text,
],
capture_output=True,
text=True,
timeout=30,
env={**os.environ, "HERMES_HOME": "/root/.hermes"},
)
if result.returncode == 0:
logger.info("Telegram send: success")
return True
logger.error("Telegram send failed (rc=%d): %s", result.returncode, result.stderr.strip()[:300])
except subprocess.TimeoutExpired:
logger.error("Telegram send: timeout")
except FileNotFoundError:
logger.error("Telegram send: hermes binary not found")
except Exception:
logger.exception("Telegram send failed")
return False
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():
"""Answer a request; optional Telegram delivery stays outside the HTTP path."""
request_started = time.perf_counter()
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 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
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
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
request_elapsed = round(time.perf_counter() - request_started, 3)
logger.info("Answer (llm=%.3fs, request=%.3fs): %s...", llm_elapsed, request_elapsed, answer[:200])
if TELEGRAM_DELIVERY_ENABLED:
_queue_telegram(answer)
return jsonify({
"status": "ok",
"answer": answer,
# 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": TELEGRAM_DELIVERY_ENABLED,
})
@app.route("/health", methods=["GET"])
async def health():
"""Return relay state. `detail=true` checks local agent state, not the LLM."""
base = {
"service": "hermes-relay",
"version": VERSION,
"mode": "async-quart",
"profile": "voice-assistant",
"telegram_delivery_enabled": TELEGRAM_DELIVERY_ENABLED,
"pending_telegram_deliveries": len(_delivery_tasks),
}
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() -> None:
global _hermes_lock
_hermes_lock = asyncio.Lock()
logger.info(
"Hermes Relay v%s startup complete (one serialized agent, telegram_delivery=%s, delivery_workers=%d, timeout=%ds)",
VERSION, TELEGRAM_DELIVERY_ENABLED, DELIVERY_WORKERS, TIMEOUT,
)
@app.after_serving
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 uvicorn
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")