File size: 7,405 Bytes
846b6be
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
"""V1 Admin alerts webhook — receives AlertManager notifications.

This is the HTTP bridge between Prometheus AlertManager and the RMI backend.
AlertManager POSTs JSON payloads here for two receivers:
  - /api/v1/admin/alerts/webhook       — default (all severities, from rmi-alerts receiver)
  - /api/v1/admin/alerts/critical      — critical only (from rmi-critical receiver)

Each handler:
  1. Parses the AlertManager v2 webhook payload
  2. Persists to Redis (sorted set, capped) for in-app display
  3. Logs structured record (JSON, single line)
  4. Returns 200 OK immediately so AlertManager doesn't retry

AlertManager webhook payload shape (Prometheus):
    {
      "version": "4",
      "groupKey": "<strings>",
      "status": "firing|resolved",
      "receiver": "rmi-alerts",
      "groupLabels": {"alertname": "...", ...},
      "commonLabels": {...},
      "commonAnnotations": {...},
      "externalURL": "...",
      "alerts": [
        {
          "status": "firing|resolved",
          "labels": {...},
          "annotations": {...},
          "startsAt": "RFC3339",
          "endsAt": "RFC3339",
          "generatorURL": "..."
        }
      ]
    }
"""
from __future__ import annotations

import json
import logging
import os
import time
from typing import Any

from fastapi import APIRouter, Request
from pydantic import BaseModel, Field

logger = logging.getLogger("rmi.admin.alerts_webhook")

# Cap the in-Redis alert history so it doesn't grow unbounded.
_REDIS_CAP = int(os.getenv("RMI_ALERTS_HISTORY_CAP", "500"))

router = APIRouter(prefix="/api/v1/admin/alerts", tags=["admin-alerts"])


class AlertmanagerAlert(BaseModel):
    status: str
    labels: dict[str, str] = Field(default_factory=dict)
    annotations: dict[str, str] = Field(default_factory=dict)
    startsAt: str | None = None
    endsAt: str | None = None
    generatorURL: str | None = None


class AlertmanagerPayload(BaseModel):
    version: str | None = None
    groupKey: str | None = None
    status: str
    receiver: str | None = None
    groupLabels: dict[str, str] = Field(default_factory=dict)
    commonLabels: dict[str, str] = Field(default_factory=dict)
    commonAnnotations: dict[str, str] = Field(default_factory=dict)
    externalURL: str | None = None
    alerts: list[AlertmanagerAlert] = Field(default_factory=list)


def _redis_url_from_env() -> str:
    """Build a Redis URL from discrete REDIS_HOST/PORT/DB/PASSWORD env vars.

    Falls back to REDIS_URL if set, otherwise localhost. Returns a URL with
    auth credentials embedded if REDIS_PASSWORD is set.
    """
    explicit = os.getenv("REDIS_URL")
    if explicit:
        return explicit
    host = os.getenv("REDIS_HOST", "localhost")
    port = os.getenv("REDIS_PORT", "6379")
    db = os.getenv("REDIS_DB", "0")
    password = os.getenv("REDIS_PASSWORD", "")
    if password:
        return f"redis://:{password}@{host}:{port}/{db}"
    return f"redis://{host}:{port}/{db}"


def _redis_client():
    """Lazy Redis import — avoids forcing a Redis dep at module load."""
    try:
        import redis.asyncio as redis_async  # type: ignore

        url = _redis_url_from_env()
        return redis_async.from_url(url, decode_responses=True)
    except Exception as exc:  # pragma: no cover - degraded mode
        logger.warning("redis_unavailable", extra={"err": str(exc)})
        return None


async def _persist_to_redis(payload: AlertmanagerPayload, severity: str) -> int:
    """Push payload to Redis sorted set capped at _REDIS_CAP entries.

    Returns number of alerts persisted (0 if Redis is down).
    """
    client = _redis_client()
    if client is None:
        return 0
    try:
        score = time.time()
        record = json.dumps(
            {
                "received_at": score,
                "severity_bucket": severity,
                "receiver": payload.receiver,
                "status": payload.status,
                "group_labels": payload.groupLabels,
                "common_labels": payload.commonLabels,
                "common_annotations": payload.commonAnnotations,
                "alerts": [a.model_dump() for a in payload.alerts],
            },
            default=str,
        )
        key = "rmi:alerts:webhook"
        async with client.pipeline(transaction=False) as pipe:
            pipe.zadd(key, {record: score})
            pipe.zremrangebyrank(key, 0, -(_REDIS_CAP + 1))
            pipe.expire(key, 7 * 24 * 3600)  # 7 days
            await pipe.execute()
        return len(payload.alerts)
    except Exception as exc:
        logger.warning("redis_persist_failed", extra={"err": str(exc)})
        return 0
    finally:
        try:
            await client.aclose()
        except Exception:
            pass


def _log_payload(payload: AlertmanagerPayload, severity: str, persisted: int) -> None:
    """Single-line JSON log so log aggregators can index cleanly."""
    record = {
        "ts": time.time(),
        "event": "alertmanager_webhook",
        "severity_bucket": severity,
        "receiver": payload.receiver,
        "status": payload.status,
        "alert_count": len(payload.alerts),
        "persisted": persisted,
        "alertname": payload.groupLabels.get("alertname"),
        "common_labels": payload.commonLabels,
    }
    logger.info(json.dumps(record, default=str))


@router.post("/webhook", status_code=200)
async def alerts_webhook(payload: AlertmanagerPayload, request: Request) -> dict[str, Any]:
    """Default receiver webhook — all severities."""
    severity = payload.commonLabels.get("severity", "unknown")
    persisted = await _persist_to_redis(payload, severity)
    _log_payload(payload, severity, persisted)
    return {
        "ok": True,
        "received": len(payload.alerts),
        "persisted": persisted,
        "severity": severity,
        "alertname": payload.groupLabels.get("alertname"),
    }


@router.post("/critical", status_code=200)
async def alerts_critical_webhook(
    payload: AlertmanagerPayload, request: Request
) -> dict[str, Any]:
    """Critical-only webhook — escalations only (severity=critical)."""
    # Defensive: if someone misroutes a non-critical here, log and accept
    # (we don't want to lose data). Severity bucket stays "critical" because
    # the receiver name implies it.
    severity = "critical"
    persisted = await _persist_to_redis(payload, severity)
    _log_payload(payload, severity, persisted)
    return {
        "ok": True,
        "received": len(payload.alerts),
        "persisted": persisted,
        "severity": severity,
        "alertname": payload.groupLabels.get("alertname"),
        "bucket": "critical",
    }


@router.get("/recent", status_code=200)
async def alerts_recent(limit: int = 50) -> dict[str, Any]:
    """Read recent alerts (debug endpoint — admin only in production)."""
    client = _redis_client()
    if client is None:
        return {"ok": False, "error": "redis_unavailable", "items": []}
    try:
        # Newest first
        raw = await client.zrevrange("rmi:alerts:webhook", 0, max(0, limit - 1))
        items = [json.loads(r) for r in raw]
        return {"ok": True, "count": len(items), "items": items}
    except Exception as exc:
        return {"ok": False, "error": str(exc), "items": []}
    finally:
        try:
            await client.aclose()
        except Exception:
            pass