Files

177 lines
5.8 KiB
Python

from __future__ import annotations
from datetime import datetime, timezone
from fastapi import FastAPI, HTTPException
from redis.asyncio import Redis
from app.config import settings
app = FastAPI(
title="alert-worker-health",
version="1.0.0",
description="Reads worker runtime status from Redis and shows queue/module health.",
)
redis_client: Redis | None = None
def _worker_status_key() -> str:
return f"{settings.redis_key_prefix}:worker:status"
def _to_bool(value: str | None) -> bool | None:
if value is None or value == "":
return None
return value.strip().lower() == "true"
def _to_int(value: str | None) -> int | None:
if value is None or value == "":
return None
try:
return int(value)
except ValueError:
return None
def _parse_iso(value: str | None) -> datetime | None:
if not value:
return None
try:
return datetime.fromisoformat(value)
except ValueError:
return None
def _is_worker_alive(heartbeat_at: str | None) -> bool:
hb = _parse_iso(heartbeat_at)
if hb is None:
return False
now = datetime.now(timezone.utc)
delta = (now - hb).total_seconds()
return delta <= max(settings.queue_block_timeout_seconds * 3, 20)
def _decode_hash(raw: dict[bytes, bytes]) -> dict[str, str]:
result: dict[str, str] = {}
for k, v in raw.items():
key = k.decode() if isinstance(k, bytes) else str(k)
value = v.decode() if isinstance(v, bytes) else str(v)
result[key] = value
return result
@app.on_event("startup")
async def startup_check() -> None:
global redis_client
if not settings.redis_enabled:
raise RuntimeError("REDIS_ENABLED must be true for worker health service")
redis_client = Redis.from_url(
settings.redis_url,
encoding="utf-8",
decode_responses=False,
)
await redis_client.ping()
@app.on_event("shutdown")
async def shutdown_event() -> None:
global redis_client
if redis_client is not None:
await redis_client.close()
@app.get("/health")
async def health() -> dict:
if redis_client is None:
raise HTTPException(status_code=503, detail="Redis client is not initialized")
key = _worker_status_key()
raw_status = await redis_client.hgetall(key)
status = _decode_hash(raw_status)
queue_len = await redis_client.llen(settings.queue_name)
processing_len = await redis_client.llen(settings.queue_processing_name)
deadletter_len = await redis_client.llen(settings.queue_deadletter_name)
heartbeat_at = status.get("heartbeat_at")
worker_alive = _is_worker_alive(heartbeat_at)
worker_state = status.get("state") or "unknown"
overall_status = "ok"
if not worker_alive:
overall_status = "degraded"
if deadletter_len > 0:
overall_status = "degraded"
if worker_state in {"error", "stopped"}:
overall_status = "degraded"
return {
"status": overall_status,
"service": "alert-worker-health",
"redis_enabled": settings.redis_enabled,
"redis_connected": True,
"queue": {
"enabled": settings.queue_enabled,
"queue_name": settings.queue_name,
"processing_name": settings.queue_processing_name,
"deadletter_name": settings.queue_deadletter_name,
"queue_length": queue_len,
"processing_length": processing_len,
"deadletter_length": deadletter_len,
},
"worker": {
"alive": worker_alive,
"state": worker_state,
"heartbeat_at": heartbeat_at,
"started_at": status.get("started_at"),
"pid": _to_int(status.get("pid")),
"current_job_id": status.get("current_job_id") or None,
"current_identity": status.get("current_identity") or None,
"current_attempt": _to_int(status.get("current_attempt")),
"last_job_id": status.get("last_job_id") or None,
"last_identity": status.get("last_identity") or None,
"last_attempt": _to_int(status.get("last_attempt")),
"last_processed_at": status.get("last_processed_at") or None,
"last_error": status.get("last_error") or None,
"last_error_at": status.get("last_error_at") or None,
"last_requeue_at": status.get("last_requeue_at") or None,
"last_deadletter_at": status.get("last_deadletter_at") or None,
"processed_count": _to_int(status.get("processed_count")) or 0,
"failed_count": _to_int(status.get("failed_count")) or 0,
"requeued_count": _to_int(status.get("requeued_count")) or 0,
"deadletter_count": _to_int(status.get("deadletter_count")) or 0,
},
"modules": {
"matrix": {
"enabled": settings.matrix_enabled,
"initialized": _to_bool(status.get("matrix_initialized")),
},
"mail": {
"enabled": settings.mail_enabled,
"initialized": _to_bool(status.get("mail_initialized")),
},
"zabbix_api": {
"enabled": settings.zabbix_api_enabled,
"initialized": _to_bool(status.get("zabbix_initialized")),
},
"llm_remediation": {
"enabled": settings.llm_enabled,
"initialized": _to_bool(status.get("llm_remediation_initialized")),
},
"llm_triage": {
"enabled": settings.llm_enabled and settings.llm_triage_enabled,
"initialized": _to_bool(status.get("llm_triage_initialized")),
},
"llm_correlation": {
"enabled": settings.llm_enabled and settings.llm_correlation_enabled,
"initialized": _to_bool(status.get("llm_correlation_initialized")),
},
},
}