From f71f5f5fc1703b054151ce9738206db82c723a00 Mon Sep 17 00:00:00 2001 From: Alexander Zubarev Date: Thu, 6 Aug 2026 18:34:46 +0300 Subject: [PATCH] =?UTF-8?q?=D0=97=D0=B0=D0=B3=D1=80=D1=83=D0=B7=D0=B8?= =?UTF-8?q?=D1=82=D1=8C=20=D1=84=D0=B0=D0=B9=D0=BB=D1=8B=20=D0=B2=20=C2=AB?= =?UTF-8?q?alert-processor/app=C2=BB?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- alert-processor/app/zabbix_client.py | 152 ++++++++++++++ alert-processor/app/zabbix_enricher.py | 270 +++++++++++++++++++++++++ 2 files changed, 422 insertions(+) create mode 100644 alert-processor/app/zabbix_client.py create mode 100644 alert-processor/app/zabbix_enricher.py diff --git a/alert-processor/app/zabbix_client.py b/alert-processor/app/zabbix_client.py new file mode 100644 index 0000000..68d9336 --- /dev/null +++ b/alert-processor/app/zabbix_client.py @@ -0,0 +1,152 @@ +from __future__ import annotations + +import itertools +from typing import Any + +import httpx + + +class ZabbixApiClient: + def __init__( + self, + api_url: str, + api_token: str, + timeout_seconds: float = 10, + verify_tls: bool = True, + ) -> None: + self.api_url = api_url + self.api_token = api_token + self._id_counter = itertools.count(1) + self.client = httpx.AsyncClient( + timeout=timeout_seconds, + verify=verify_tls, + headers={ + "Authorization": f"Bearer {self.api_token}", + "Content-Type": "application/json-rpc", + }, + ) + + async def close(self) -> None: + await self.client.aclose() + + async def _call(self, method: str, params: dict[str, Any]) -> Any: + payload = { + "jsonrpc": "2.0", + "method": method, + "params": params, + "id": next(self._id_counter), + } + + response = await self.client.post(self.api_url, json=payload) + response.raise_for_status() + + data = response.json() + if "error" in data: + error = data["error"] + raise RuntimeError( + f"Zabbix API error: method={method} " + f"code={error.get('code')} message={error.get('message')} " + f"data={error.get('data')}" + ) + + return data.get("result") + + async def get_trigger(self, trigger_id: str) -> dict[str, Any] | None: + result = await self._call( + "trigger.get", + { + "output": [ + "triggerid", + "description", + "comments", + "opdata", + "priority", + "state", + "status", + "value", + ], + "triggerids": [trigger_id], + "selectHosts": ["hostid", "host", "name"], + "selectItems": ["itemid", "name", "key_", "lastvalue", "units", "value_type"], + "selectTags": "extend", + "limit": 1, + }, + ) + return result[0] if result else None + + async def get_item(self, item_id: str) -> dict[str, Any] | None: + result = await self._call( + "item.get", + { + "output": [ + "itemid", + "name", + "key_", + "lastvalue", + "units", + "value_type", + "status", + ], + "itemids": [item_id], + "selectHosts": ["hostid", "host", "name"], + "selectTags": "extend", + "limit": 1, + }, + ) + return result[0] if result else None + + async def get_event(self, event_id: str) -> dict[str, Any] | None: + result = await self._call( + "event.get", + { + "output": [ + "eventid", + "objectid", + "clock", + "name", + "severity", + "value", + "acknowledged", + ], + "eventids": [event_id], + "selectHosts": ["hostid", "host", "name"], + "selectTags": "extend", + "limit": 1, + }, + ) + return result[0] if result else None + + async def get_host(self, host_id: str) -> dict[str, Any] | None: + result = await self._call( + "host.get", + { + "output": ["hostid", "host", "name"], + "hostids": [host_id], + "selectTags": "extend", + "limit": 1, + }, + ) + return result[0] if result else None + + async def get_history( + self, + item_id: str, + value_type: int, + time_from: int, + time_till: int, + limit: int = 500, + ) -> list[dict[str, Any]]: + result = await self._call( + "history.get", + { + "output": ["clock", "value"], + "history": value_type, + "itemids": [item_id], + "time_from": time_from, + "time_till": time_till, + "sortfield": "clock", + "sortorder": "ASC", + "limit": limit, + }, + ) + return result or [] diff --git a/alert-processor/app/zabbix_enricher.py b/alert-processor/app/zabbix_enricher.py new file mode 100644 index 0000000..c89c487 --- /dev/null +++ b/alert-processor/app/zabbix_enricher.py @@ -0,0 +1,270 @@ +from __future__ import annotations + +import re +from datetime import datetime, timedelta +from pathlib import Path +from urllib.parse import urlencode +from zoneinfo import ZoneInfo + +from app.models import NormalizedEvent +from app.tag_utils import build_correlation_scopes, infer_domain_from_text, merge_tags, normalize_tag_dict +from app.zabbix_client import ZabbixApiClient + + +def _contains_unresolved_macros(value: str | None) -> bool: + if not value: + return False + return bool(re.search(r"\{[A-Z0-9\._:#\$]+\}", value)) + + +def _safe_filename(value: str) -> str: + return re.sub(r"[^a-zA-Z0-9._-]+", "_", value).strip("_") or "graph" + + +def _detect_alert_scope( + item_key: str | None, + service: str | None, + trigger_name: str | None, + tags: dict[str, str] | None = None, +) -> str: + item_key_lower = (item_key or "").strip().lower() + service_lower = (service or "").strip().lower() + trigger_lower = (trigger_name or "").strip().lower() + tags = tags or {} + + if tags.get("scope"): + return tags["scope"] + if tags.get("service"): + return tags["service"] + + tag_values = " ".join(f"{k}={v}" for k, v in tags.items()).lower() + + host_prefixes = ( + "system.cpu.", + "system.swap.", + "system.uptime", + "vm.memory.", + "vfs.fs.", + "vfs.dev.", + "proc.num", + "kernel.", + "net.if.", + "agent.", + ) + + container_prefixes = ( + "docker.", + "container.", + "podman.", + "kube.", + "kubernetes.", + "cri.", + ) + + if item_key_lower.startswith(host_prefixes): + return "host_os" + + if item_key_lower.startswith(container_prefixes): + return "container" + + host_hints = ( + "cpu utilization", + "memory utilization", + "filesystem", + "disk space", + "load average", + "linux:", + "windows:", + ) + + container_hints = ( + "container", + "docker", + "pod", + "kube", + "kubernetes", + ) + + combined = " ".join([service_lower, trigger_lower, tag_values]) + + if any(h in combined for h in container_hints): + return "container" + + if any(h in combined for h in host_hints): + return "host_os" + + return "unknown" + + +class ZabbixEnricher: + def __init__( + self, + client: ZabbixApiClient, + web_url: str | None = None, + graph_period_hours: int = 1, + graph_image_dir: str = ".graph_images", + graph_timezone: str = "UTC", + ) -> None: + self.client = client + self.web_url = (web_url or "").rstrip("/") + self.graph_period_hours = graph_period_hours + self.graph_image_dir = Path(graph_image_dir) + self.graph_image_dir.mkdir(parents=True, exist_ok=True) + self.graph_timezone_name = graph_timezone or "UTC" + self.graph_timezone = ZoneInfo(self.graph_timezone_name) + + async def enrich_event(self, event: NormalizedEvent) -> NormalizedEvent: + context: dict[str, object] = {} + + if self.web_url and event.trigger_id and event.event_id: + event.event_url = ( + f"{self.web_url}/tr_events.php" + f"?triggerid={event.trigger_id}&eventid={event.event_id}" + ) + elif event.zabbix_url and not event.event_url: + event.event_url = event.zabbix_url + + item_value_type: int | None = None + item_name: str | None = None + item_units: str | None = None + item_key: str | None = None + + incoming_tags = normalize_tag_dict(event.tags) + payload_tags = normalize_tag_dict(event.raw_payload.get("tags")) + + trigger_tags: dict[str, str] = {} + host_tags: dict[str, str] = {} + item_tags: dict[str, str] = {} + event_tags: dict[str, str] = {} + + if event.trigger_id: + trigger = await self.client.get_trigger(event.trigger_id) + if trigger: + trigger_tags = normalize_tag_dict(trigger.get("tags")) + hosts = trigger.get("hosts") or [] + if hosts and not event.host: + event.host = hosts[0].get("host") or hosts[0].get("name") or event.host + if hosts and not event.host_id: + event.host_id = hosts[0].get("hostid") or event.host_id + + items = trigger.get("items") or [] + if items and not event.item_id: + event.item_id = items[0].get("itemid") or event.item_id + + description = trigger.get("description") + comments = trigger.get("comments") + trigger_opdata = trigger.get("opdata") + if description: + context["trigger_description"] = str(description) + if comments: + context["trigger_comments"] = str(comments) + + if trigger_opdata and not _contains_unresolved_macros(str(trigger_opdata)): + context["resolved_opdata"] = str(trigger_opdata) + + if event.host_id: + host_info = await self.client.get_host(event.host_id) + if host_info: + host_tags = normalize_tag_dict(host_info.get("tags")) + + if event.item_id: + item = await self.client.get_item(event.item_id) + if item: + item_tags = normalize_tag_dict(item.get("tags")) + + if not event.service: + event.service = item.get("name") or event.service + + item_name = item.get("name") + item_key = item.get("key_") + lastvalue = item.get("lastvalue") + item_units = item.get("units") or "" + try: + item_value_type = int(item.get("value_type")) + except (TypeError, ValueError): + item_value_type = None + + if item_name: + context["item_name"] = str(item_name) + if item_key: + context["item_key"] = str(item_key) + if lastvalue not in (None, ""): + context["last_value"] = f"{lastvalue}{(' ' + item_units) if item_units else ''}" + + if event.event_id: + zbx_event = await self.client.get_event(event.event_id) + if zbx_event: + event_tags = normalize_tag_dict(zbx_event.get("tags")) + event_name = zbx_event.get("name") + if event_name and not event.trigger_name: + event.trigger_name = str(event_name) + + merged_tags = merge_tags(host_tags, trigger_tags, item_tags, event_tags, payload_tags, incoming_tags) + + if "service" not in merged_tags and event.service: + merged_tags["service"] = str(event.service) + + domain = merged_tags.get("domain") or infer_domain_from_text( + event.trigger_name, + event.zabbix_url, + event.event_url, + event.raw_payload.get("url"), + event.raw_payload.get("name"), + ) + if domain: + merged_tags["domain"] = domain + + alert_scope = _detect_alert_scope( + item_key=item_key, + service=merged_tags.get("service") or event.service, + trigger_name=event.trigger_name, + tags=merged_tags, + ) + + correlation_scopes = build_correlation_scopes( + tags=merged_tags, + host=event.host, + trigger_name=event.trigger_name, + ) + + if merged_tags: + event.tags = merged_tags + + if not context.get("resolved_opdata"): + if event.opdata and not _contains_unresolved_macros(event.opdata): + context["resolved_opdata"] = event.opdata + elif context.get("last_value"): + metric_label = context.get("item_name") or event.service or "Metric value" + context["resolved_opdata"] = f"{metric_label}: {context['last_value']}" + + if self.web_url and event.item_id: + now_local = datetime.now(self.graph_timezone) + from_local = now_local - timedelta(hours=self.graph_period_hours) + + params = { + "from": from_local.strftime("%Y-%m-%d %H:%M:%S"), + "to": now_local.strftime("%Y-%m-%d %H:%M:%S"), + "itemids[0]": event.item_id, + } + event.graph_url = f"{self.web_url}/chart.php?{urlencode(params)}" + + if item_value_type is not None and item_value_type in {0, 3}: + file_name = _safe_filename(f"{event.host}_{event.item_id}_{event.event_id or 'graph'}.png") + event.graph_image_path = str(self.graph_image_dir / file_name) + else: + event.graph_image_path = None + + context["alert_scope"] = alert_scope + context["service"] = merged_tags.get("service") + context["scope"] = merged_tags.get("scope") + context["component"] = merged_tags.get("component") + context["domain"] = merged_tags.get("domain") + context["tags"] = merged_tags + context["host_tags"] = host_tags + context["trigger_tags"] = trigger_tags + context["item_tags"] = item_tags + context["event_tags"] = event_tags + context["correlation_scopes"] = correlation_scopes + + event.zabbix_context.update(context) + return event