diff --git a/alert-receiver/app/normalizer.py b/alert-receiver/app/normalizer.py new file mode 100644 index 0000000..4ac6ad1 --- /dev/null +++ b/alert-receiver/app/normalizer.py @@ -0,0 +1,152 @@ +from __future__ import annotations + +import json +from datetime import datetime, timezone +from typing import Any + + +SEVERITY_MAP = { + 0: "Not classified", + 1: "Information", + 2: "Warning", + 3: "Average", + 4: "High", + 5: "Disaster", +} + + +def _first_non_empty(payload: dict[str, Any], keys: list[str]) -> Any: + for key in keys: + value = payload.get(key) + if value is not None and str(value).strip() != "": + return value + return None + + +def _parse_timestamp(value: Any) -> datetime | None: + if value is None: + return None + + if isinstance(value, (int, float)): + return datetime.fromtimestamp(value, tz=timezone.utc) + + if isinstance(value, str): + raw = value.strip() + if not raw: + return None + + if raw.isdigit(): + return datetime.fromtimestamp(int(raw), tz=timezone.utc) + + try: + return datetime.fromisoformat(raw.replace("Z", "+00:00")) + except ValueError: + return None + + return None + + +def _parse_severity(payload: dict[str, Any]) -> tuple[str | None, int | None]: + raw = _first_non_empty( + payload, + [ + "severity", + "severity_name", + "event_severity", + "event_severity_name", + "priority", + ], + ) + + if raw is None: + return None, None + + if isinstance(raw, int): + return SEVERITY_MAP.get(raw), raw + + raw_str = str(raw).strip() + if raw_str.isdigit(): + code = int(raw_str) + return SEVERITY_MAP.get(code, raw_str), code + + return raw_str, None + + +def _parse_tags(raw_tags: Any) -> dict[str, str]: + if raw_tags is None: + return {} + + if isinstance(raw_tags, dict): + return {str(k): str(v) for k, v in raw_tags.items()} + + if isinstance(raw_tags, list): + result: dict[str, str] = {} + for item in raw_tags: + if isinstance(item, dict): + tag = item.get("tag") or item.get("name") or item.get("key") + value = item.get("value") + if tag is not None: + result[str(tag)] = "" if value is None else str(value) + return result + + if isinstance(raw_tags, str): + raw = raw_tags.strip() + if not raw: + return {} + try: + parsed = json.loads(raw) + return _parse_tags(parsed) + except json.JSONDecodeError: + return {} + + return {} + + +def normalize_zabbix_payload( + payload: dict[str, Any], + correlation_id: str, + remote_addr: str | None, +) -> dict[str, Any]: + severity, severity_code = _parse_severity(payload) + + timestamp = _parse_timestamp( + _first_non_empty(payload, ["timestamp", "clock", "event_time"]) + ) + + tags = _parse_tags(_first_non_empty(payload, ["tags", "event_tags"])) + + trigger_name = _first_non_empty( + payload, + ["trigger_name", "name", "event_name", "problem_name", "subject"], + ) + + service = _first_non_empty( + payload, + ["service", "component", "application", "item_name"], + ) + + normalized = { + "source": "zabbix", + "correlation_id": correlation_id, + "received_at": datetime.now(tz=timezone.utc), + "remote_addr": remote_addr, + "event_id": _first_non_empty(payload, ["event_id", "eventid"]), + "problem_id": _first_non_empty(payload, ["problem_id", "problemid"]), + "event_type": _first_non_empty(payload, ["event_type", "status"]) or "problem", + "timestamp": timestamp, + "severity": severity, + "severity_code": severity_code, + "host": _first_non_empty(payload, ["host", "host_name", "hostname"]), + "host_id": _first_non_empty(payload, ["host_id", "hostid"]), + "service": service, + "trigger_name": trigger_name, + "trigger_id": _first_non_empty(payload, ["trigger_id", "triggerid"]), + "item_id": _first_non_empty(payload, ["item_id", "itemid"]), + "value": _first_non_empty(payload, ["value", "status_value"]), + "opdata": _first_non_empty(payload, ["opdata", "message", "description"]), + "zabbix_url": _first_non_empty(payload, ["zabbix_url", "url", "event_url"]), + "tags": tags, + "raw_payload": payload, + } + + return normalized \ No newline at end of file