Files

531 lines
18 KiB
Python

from __future__ import annotations
import logging
import re
from dataclasses import dataclass, field
from pathlib import Path
from typing import Any
import yaml
from app.config import settings
from app.models import NotificationDecision, ProcessorForwardEnvelope
logger = logging.getLogger(__name__)
def _norm(value: str | None) -> str:
return (value or "").strip().lower()
def severity_rank(severity: str | None) -> int:
value = _norm(severity)
if value in {"disaster"}:
return 4
if value in {"high"}:
return 3
if value in {"average"}:
return 2
if value in {"warning", "information", "not classified", "not_classified"}:
return 1
return 0
@dataclass(frozen=True)
class EventKindRule:
kind: str
trigger_patterns: tuple[str, ...] = ()
item_key_patterns: tuple[str, ...] = ()
service_patterns: tuple[str, ...] = ()
tag_patterns: tuple[str, ...] = ()
scope_in: tuple[str, ...] = ()
severity_in: tuple[str, ...] = ()
def matches(
self,
envelope: ProcessorForwardEnvelope,
decision: NotificationDecision,
) -> bool:
event = envelope.event
zbx_context = event.zabbix_context or {}
trigger_name = _norm(event.trigger_name)
item_key = _norm(zbx_context.get("item_key"))
service = _norm(zbx_context.get("service") or event.service)
alert_scope = _norm(zbx_context.get("alert_scope"))
severity = _norm(decision.severity)
tags_joined = " ".join(
f"{_norm(str(k))}={_norm(str(v))}" for k, v in (event.tags or {}).items()
)
if self.scope_in and alert_scope not in self.scope_in:
return False
if self.severity_in and severity not in self.severity_in:
return False
if self.trigger_patterns and not any(re.search(pattern, trigger_name) for pattern in self.trigger_patterns):
return False
if self.item_key_patterns and not any(re.search(pattern, item_key) for pattern in self.item_key_patterns):
return False
if self.service_patterns and not any(re.search(pattern, service) for pattern in self.service_patterns):
return False
if self.tag_patterns and not all(re.search(pattern, tags_joined) for pattern in self.tag_patterns):
return False
return True
@dataclass
class CorrelationEventRecord:
host: str
event_id: str | None
correlation_id: str
kind: str
severity: str | None
routing_class: str | None
fingerprint: str | None
root_candidate: bool
role: str
group_id: str | None
parent_event_id: str | None
parent_correlation_id: str | None
timestamp: float
tags: dict[str, str] = field(default_factory=dict)
service: str | None = None
scope: str | None = None
component: str | None = None
domain: str | None = None
@dataclass
class CorrelationAssessment:
applied: bool
role: str
kind: str
group_id: str | None
reason: str | None
parent_event_id: str | None = None
parent_correlation_id: str | None = None
root_cause_candidate: bool = False
correlated_event_count: int = 0
suppress_child: bool = False
source: str = "deterministic"
confidence: str | None = "high"
@dataclass
class CorrelationRegistry:
event_kind_rules: list[EventKindRule] = field(default_factory=list)
root_cause_map: dict[str, set[str]] = field(default_factory=dict)
def classify_event_kind(
self,
envelope: ProcessorForwardEnvelope,
decision: NotificationDecision,
) -> str:
for rule in self.event_kind_rules:
if rule.matches(envelope, decision):
return rule.kind
return "unknown"
def is_root_candidate(self, kind: str) -> bool:
return kind in self.root_cause_map
def explains(self, parent_kind: str, child_kind: str) -> bool:
if parent_kind == child_kind:
return False
return child_kind in self.root_cause_map.get(parent_kind, set())
@classmethod
def default(cls) -> "CorrelationRegistry":
return cls(
event_kind_rules=[
EventKindRule(
kind="postgresql_unavailable",
trigger_patterns=(r"postgresql.+unavailable", r"postgres.+unavailable"),
),
EventKindRule(
kind="host_down",
trigger_patterns=(r"host unavailable", r"icmp.+unreachable", r"unreachable"),
scope_in=("host_os",),
),
EventKindRule(
kind="network_port_down",
trigger_patterns=(r"port down", r"interface down"),
),
EventKindRule(
kind="ap_down",
trigger_patterns=(r"ap down",),
),
EventKindRule(
kind="disk_full",
item_key_patterns=(r"^vfs\.fs\.",),
trigger_patterns=(r"space", r"full", r"disk"),
scope_in=("host_os",),
),
EventKindRule(
kind="container_unavailable",
trigger_patterns=(
r"container",
r"docker",
r"health state",
r"not running",
r"stopped",
r"unavailable",
r"down",
r"failed",
),
scope_in=("container",),
),
EventKindRule(
kind="cpu_high",
item_key_patterns=(r"^system\.cpu\.",),
trigger_patterns=(r"cpu utilization",),
scope_in=("host_os",),
),
EventKindRule(
kind="memory_high",
item_key_patterns=(r"^vm\.memory\.", r"^system\.swap\."),
trigger_patterns=(r"memory utilization",),
scope_in=("host_os",),
),
EventKindRule(
kind="ssl_expiry",
trigger_patterns=(r"ssl", r"certificate", r"expire"),
),
EventKindRule(
kind="service_unavailable",
trigger_patterns=(r"service unavailable", r"unavailable"),
),
EventKindRule(
kind="site_unavailable",
trigger_patterns=(r"site.+is down", r"failed step of scenario"),
),
EventKindRule(
kind="db_connection_error",
trigger_patterns=(r"connection", r"db"),
),
],
root_cause_map={
"host_down": {
"service_unavailable",
"container_unavailable",
"postgresql_unavailable",
"cpu_high",
"memory_high",
"disk_full",
"ssl_expiry",
"db_connection_error",
"site_unavailable",
},
"network_port_down": {
"host_down",
"service_unavailable",
"ap_down",
"site_unavailable",
},
"postgresql_unavailable": {
"service_unavailable",
"db_connection_error",
"site_unavailable",
},
"disk_full": {
"postgresql_unavailable",
"container_unavailable",
"service_unavailable",
"db_connection_error",
"site_unavailable",
},
"container_unavailable": {
"service_unavailable",
"site_unavailable",
},
},
)
@classmethod
def from_yaml_files(
cls,
event_kind_rules_path: str,
root_cause_map_path: str,
) -> "CorrelationRegistry":
base = cls.default()
custom_kind_rules = _load_event_kind_rules(event_kind_rules_path)
custom_root_cause_map = _load_root_cause_map(root_cause_map_path)
merged_kind_rules = custom_kind_rules + base.event_kind_rules
merged_root_cause_map: dict[str, set[str]] = {
root_kind: set(children)
for root_kind, children in base.root_cause_map.items()
}
for root_kind, explained in custom_root_cause_map.items():
if root_kind not in merged_root_cause_map:
merged_root_cause_map[root_kind] = set()
merged_root_cause_map[root_kind].update(explained)
logger.info(
"Correlation registry loaded: custom_kind_rules=%s custom_root_kinds=%s total_kind_rules=%s total_root_kinds=%s",
len(custom_kind_rules),
len(custom_root_cause_map),
len(merged_kind_rules),
len(merged_root_cause_map),
)
return cls(
event_kind_rules=merged_kind_rules,
root_cause_map=merged_root_cause_map,
)
def _load_yaml(path: str) -> dict[str, Any]:
file_path = Path(path)
if not file_path.exists():
logger.warning("Correlation YAML file not found: %s. Using defaults/merge fallback.", path)
return {}
try:
with file_path.open("r", encoding="utf-8") as f:
data = yaml.safe_load(f) or {}
if not isinstance(data, dict):
logger.warning("Correlation YAML root must be a mapping: %s", path)
return {}
return data
except Exception as exc:
logger.exception("Failed to load correlation YAML %s: %s", path, exc)
return {}
def _load_event_kind_rules(path: str) -> list[EventKindRule]:
data = _load_yaml(path)
raw_rules = data.get("event_kind_rules", [])
if not isinstance(raw_rules, list):
logger.warning("event_kind_rules must be a list in %s", path)
return []
rules: list[EventKindRule] = []
for idx, item in enumerate(raw_rules, start=1):
if not isinstance(item, dict):
logger.warning("Skipping invalid event_kind_rules[%s] in %s", idx, path)
continue
kind = _norm(item.get("kind"))
if not kind:
logger.warning("Skipping event_kind_rules[%s] without kind in %s", idx, path)
continue
rules.append(
EventKindRule(
kind=kind,
trigger_patterns=tuple(item.get("trigger_patterns", []) or []),
item_key_patterns=tuple(item.get("item_key_patterns", []) or []),
service_patterns=tuple(item.get("service_patterns", []) or []),
tag_patterns=tuple(item.get("tag_patterns", []) or []),
scope_in=tuple(_norm(v) for v in (item.get("scope_in", []) or [])),
severity_in=tuple(_norm(v) for v in (item.get("severity_in", []) or [])),
)
)
return rules
def _load_root_cause_map(path: str) -> dict[str, set[str]]:
data = _load_yaml(path)
raw_map = data.get("root_cause_map", {})
if not isinstance(raw_map, dict):
logger.warning("root_cause_map must be a mapping in %s", path)
return {}
result: dict[str, set[str]] = {}
for root_kind_raw, payload in raw_map.items():
root_kind = _norm(root_kind_raw)
if not root_kind:
continue
explains: list[str] = []
if isinstance(payload, dict):
explains_raw = payload.get("explains", [])
if isinstance(explains_raw, list):
explains = [_norm(v) for v in explains_raw if _norm(v)]
elif isinstance(payload, list):
explains = [_norm(v) for v in payload if _norm(v)]
if explains:
result[root_kind] = set(explains)
return result
def _shared_service_scope_or_domain(
envelope: ProcessorForwardEnvelope,
candidate: CorrelationEventRecord,
) -> bool:
event = envelope.event
zbx_context = event.zabbix_context or {}
current_tags = {str(k).lower(): str(v).strip() for k, v in (event.tags or {}).items()}
current_service = _norm(current_tags.get("service") or zbx_context.get("service") or event.service)
current_scope = _norm(current_tags.get("scope") or zbx_context.get("scope"))
current_domain = _norm(current_tags.get("domain") or zbx_context.get("domain"))
candidate_service = _norm(candidate.service or candidate.tags.get("service"))
candidate_scope = _norm(candidate.scope or candidate.tags.get("scope"))
candidate_domain = _norm(candidate.domain or candidate.tags.get("domain"))
if current_scope and candidate_scope and current_scope == candidate_scope:
return True
if current_service and candidate_service and current_service == candidate_service:
return True
if current_domain and candidate_domain and current_domain == candidate_domain:
return True
if event.host and candidate.host and _norm(event.host) == _norm(candidate.host):
return True
return False
def assess_correlation(
envelope: ProcessorForwardEnvelope,
decision: NotificationDecision,
recent_events: list[CorrelationEventRecord],
registry: CorrelationRegistry,
) -> CorrelationAssessment:
kind = registry.classify_event_kind(envelope, decision)
root_candidate = registry.is_root_candidate(kind)
if kind == "unknown":
return CorrelationAssessment(
applied=False,
role="standalone",
kind=kind,
group_id=None,
reason=None,
root_cause_candidate=False,
correlated_event_count=0,
)
parent_candidates = [
item
for item in recent_events
if item.root_candidate
and registry.explains(item.kind, kind)
and _shared_service_scope_or_domain(envelope, item)
]
parent_candidates.sort(
key=lambda item: (item.timestamp, severity_rank(item.severity)),
reverse=True,
)
if parent_candidates:
parent = parent_candidates[0]
suppress_child = (
settings.correlation_suppress_children
and decision.routing_class in {
"average_priority",
"low_priority",
"triage_low_priority_notify",
"triage_low_priority_hold",
"triage_low_priority_suppressed",
}
)
return CorrelationAssessment(
applied=True,
role="child",
kind=kind,
group_id=parent.group_id or parent.event_id or parent.correlation_id,
reason=(
f"Likely downstream of {parent.kind}; matched recent event "
f"in the same service/scope/domain window within {settings.correlation_window_seconds}s"
),
parent_event_id=parent.event_id,
parent_correlation_id=parent.correlation_id,
root_cause_candidate=False,
correlated_event_count=0,
suppress_child=suppress_child,
)
if root_candidate:
correlated_count = sum(
1
for item in recent_events
if registry.explains(kind, item.kind) and _shared_service_scope_or_domain(envelope, item)
)
return CorrelationAssessment(
applied=True,
role="root",
kind=kind,
group_id=envelope.event.event_id or envelope.event.correlation_id,
reason=(
f"Potential root cause candidate; {correlated_count} related alerts seen "
f"in matching service/scope/domain window within {settings.correlation_window_seconds}s"
),
root_cause_candidate=True,
correlated_event_count=correlated_count,
suppress_child=False,
)
return CorrelationAssessment(
applied=True,
role="standalone",
kind=kind,
group_id=envelope.event.event_id or envelope.event.correlation_id,
reason="No root cause candidate found in matching service/scope/domain correlation window",
root_cause_candidate=False,
correlated_event_count=0,
suppress_child=False,
)
def apply_correlation_to_decision(
decision: NotificationDecision,
assessment: CorrelationAssessment,
) -> NotificationDecision:
if not assessment.applied:
return decision
update = {
"correlation_applied": True,
"correlation_role": assessment.role,
"correlation_group_id": assessment.group_id,
"correlation_kind": assessment.kind,
"correlation_reason": assessment.reason,
"correlation_source": assessment.source,
"correlation_confidence": assessment.confidence,
"parent_event_id": assessment.parent_event_id,
"parent_correlation_id": assessment.parent_correlation_id,
"root_cause_candidate": assessment.root_cause_candidate,
"correlated_event_count": assessment.correlated_event_count,
}
if assessment.role == "child" and assessment.suppress_child:
update.update(
{
"notify": False,
"suppressed": True,
"channels": [],
"routing_class": "correlated_child_suppressed",
"reason": "Correlated child alert suppressed by RCA baseline",
"suppress_reason": assessment.reason,
}
)
return decision.model_copy(update=update)
if assessment.role == "child":
update["reason"] = f"{decision.reason}. RCA baseline matched a parent event in the same service/scope."
return decision.model_copy(update=update)
if assessment.role == "root":
update["reason"] = f"{decision.reason}. RCA baseline marked this event as root cause candidate"
return decision.model_copy(update=update)
return decision.model_copy(update=update)