alert_engine.py
"""
AlertEngine: motor de deduplicación de registros en streaming y limitación de velocidad de alertas.
Librería estándar pura, sin hilos en segundo plano. Utiliza tiempos 'now' suministrados externamente.
"""
from future import annotations
import re
import hashlib
from collections import deque, defaultdict
from dataclasses import dataclass
from typing import Dict, Tuple, Optional, List, Any
Mapeo de severidad
_SEVERITY_ORDER = {
"debug": 10,
"info": 20,
"warn": 30,
"warning": 30, # acepta alternativa
"error": 40,
"critical": 50,
}
MAX_SAMPLE_LEN = 200
MAX_FINGERPRINTS = 100000 # límite de memoria determinista; ver nota de diseño
_uuid_re = re.compile(r"\b[0-9a-fA-F]{8}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{12}\b")
_ipv4_re = re.compile(r"\b(?:\d{1,3}.){3}\d{1,3}\b")
_hexid_re = re.compile(r"\b[0-9a-fA-F]{8,}\b")
_digits_re = re.compile(r"\d+")
def _normalize_message(msg: str) -> str:
# Reemplaza UUIDs, IPs, ids hexadecimales largos, luego secuencias de dígitos. El orden importa.
s = _uuid_re.sub("<UUID>", msg)
s = _ipv4_re.sub("<IP>", s)
s = _hexid_re.sub("<HEX>", s)
s = _digits_re.sub("<NUM>", s)
return s
def _fingerprint_for(service: str, severity: str, message: str) -> Tuple[str, str]:
normalized = _normalize_message(message)
key_str = f"{service}|{severity}|{normalized}"
fp = hashlib.sha1(key_str.encode("utf-8", errors="ignore")).hexdigest()
return fp, normalized
@dataclass
class FPEntry:
service: str
severity: str
normalized: str
first_seen: float
last_seen: float
count: int
sample_message: str
suppressed_since_emit: int
burst_times: deque # marcas de tiempo (flotantes) de duplicados suprimidos para detección de ráfagas
last_emitted: Optional[float]
class AlertEngine:
"""AlertEngine(config: dict, now: float)
claves de config:
- window_seconds: float
- max_alerts_per_window: int
- dedup_seconds: float
- severity_floor: uno de debug, info, warn, error, critical
- burst_escalation: dict opcional {"count": int, "within_seconds": float}
Política para 'now' no monótono: el motor mantiene una marca de agua interna (el 'now' máximo visto).
Si ingest() o flush() se llama con un 'now' anterior a la marca de agua, el motor trata el 'now' proporcionado
como la marca de agua (es decir, ajusta el tiempo hacia adelante).
Esto asegura que las alertas emitidas nunca se produzcan fuera de orden por tiempo y hace que el
comportamiento sea determinista para eventos tardíos.
Política de desalojo: para limitar la memoria, se aplica un límite determinista MAX_FINGERPRINTS.
Cuando se excede, el motor desalojará las huellas dactilares con la más antigua (menor) last_seen,
rompiendo empates por first_seen. Esto es determinista y evita el crecimiento ilimitado.
"""
def __init__(self, config: dict, now: float):
# validación mínima de config
try:
self.window_seconds = float(config["window_seconds"])
self.max_alerts_per_window = int(config["max_alerts_per_window"])
self.dedup_seconds = float(config["dedup_seconds"])
self.severity_floor = str(config["severity_floor"]).lower()
if self.severity_floor not in _SEVERITY_ORDER:
raise KeyError
except Exception:
raise ValueError("configuración inválida: claves requeridas window_seconds, max_alerts_per_window, dedup_seconds, severity_floor")
be = config.get("burst_escalation")
if be is not None:
try:
self.burst_count = int(be["count"])
self.burst_within = float(be["within_seconds"])
if self.burst_count <= 0 or self.burst_within <= 0:
raise ValueError
except Exception:
raise ValueError("configuración de burst_escalation inválida")
else:
self.burst_count = None
self.burst_within = None
# estado interno
self._now_watermark = float(now)
self._fp_store: Dict[str, FPEntry] = {}
# cola por clave (servicio, severidad) de marcas de tiempo de emisión para ventana deslizante
self._key_emissions: Dict[Tuple[str, str], deque] = defaultdict(deque)
# estadísticas
self._total_ingested = 0
self._total_emitted = 0
self._total_suppressed = 0
def _clamp_now(self, now: float) -> float:
if now < self._now_watermark:
# política: tratar como tardío; ajustar hacia adelante a la marca de agua
return self._now_watermark
self._now_watermark = now
return now
def _ensure_eviction(self) -> None:
if len(self._fp_store) <= MAX_FINGERPRINTS:
return
# desalojo determinista: ordenar por last_seen, luego first_seen
items = sorted(self._fp_store.items(), key=lambda kv: (kv[1].last_seen, kv[1].first_seen))
to_evict = len(self._fp_store) - MAX_FINGERPRINTS
for i in range(to_evict):
k = items[i][0]
del self._fp_store[k]
def ingest(self, event: dict, now: float) -> List[dict]:
"""Procesa un solo evento y devuelve una lista de alertas emitidas (puede estar vacía).
Política para entrada mal formada: el motor cuenta el evento como ingerido pero lo descarta
y devuelve una lista vacía cuando faltan campos requeridos o la severidad es inválida.
"""
self._total_ingested += 1
now = self._clamp_now(float(now))
alerts: List[dict] = []
# validación mínima
if not isinstance(event, dict):
return []
required = ("timestamp", "severity", "service", "message")
for r in required:
if r not in event:
return []
try:
ev_ts = float(event["timestamp"])
severity = str(event["severity"]).lower()
service = str(event["service"])
message = str(event["message"])
labels = event.get("labels")
if labels is not None and not isinstance(labels, dict):
labels = None
except Exception:
return []
if severity not in _SEVERITY_ORDER:
# manejo explícito: severidad desconocida -> descarta evento
return []
# umbral de severidad
if _SEVERITY_ORDER[severity] < _SEVERITY_ORDER[self.severity_floor]:
# descartado antes de cualquier otro procesamiento
return []
# normalizar y crear huella digital
fp, normalized = _fingerprint_for(service, severity, message)
# asegurar que sample_message se trunque sin afectar la huella digital
sample_message = message[:MAX_SAMPLE_LEN]
entry = self._fp_store.get(fp)
if entry is None:
entry = FPEntry(
service=service,
severity=severity,
normalized=normalized,
first_seen=ev_ts,
last_seen=ev_ts,
count=1,
sample_message=sample_message,
suppressed_since_emit=0,
burst_times=deque(),
last_emitted=None,
)
self._fp_store[fp] = entry
else:
# actualizar marcas de tiempo/contadores
entry.count += 1
entry.last_seen = ev_ts
# mantener la primera sample_message
if len(entry.sample_message) < MAX_SAMPLE_LEN:
# intentar completar la muestra si la inicial fue corta
entry.sample_message = (entry.sample_message + " | " + sample_message)[:MAX_SAMPLE_LEN]
# verificación de deduplicación: si está dentro de dedup_seconds de last_seen (usando la marca de tiempo del evento), tratar como duplicado
is_duplicate = (ev_ts - entry.last_seen) <= self.dedup_seconds if entry.count > 1 else False
# Nota: last_seen ya se ha establecido en ev_ts; para la lógica de duplicados deberíamos usar el last_seen anterior
# Para implementar correctamente, recalcular: si count==1 -> no es duplicado; si no, si ev_ts - prev_last_seen <= dedup_seconds
if entry.count == 1:
is_duplicate = False
else:
# necesitamos el last_seen anterior: aproximar comprobando ev_ts - (entry.last_seen o ev_ts) pero actualizamos last_seen.
# Más simple: consideramos duplicados como eventos cuya marca de tiempo del evento está dentro de dedup_seconds de la marca de tiempo del evento anterior.
# Como actualizamos, no podemos recuperar el anterior; corregir manteniendo temporalmente last_seen_old antes.
pass
# La lógica anterior es incómoda porque actualizamos last_seen prematuramente. Para corregir, rehacer con seguimiento previo adecuado.
# Recalcular con almacenamiento explícito.
# Reconstruir: obtener instantánea de entrada previa
prev_last_seen = entry.last_seen
prev_count = entry.count
# Recomputar correctamente: si prev_count > 1, el last_seen anterior no es accesible; para evitar complejidad,
# trataremos los duplicados usando now (tiempo de procesamiento) en relación con entry.last_seen almacenado antes del evento actual.
# Para lograr esto, almacenamos last_seen igual a ev_ts arriba; pero si prev_count>=1, el last_seen anterior es igual a stored_entry.first_seen para count==1,
# o algo más. Esta complejidad surge porque actualizamos campos prematuramente.
# Para simplificar y hacer determinista: consideraremos duplicados si now - entry.last_seen <= dedup_seconds usando last_seen antes de actualizar.
# Como no podemos recuperar el last_seen anterior, en su lugar mantendremos la detección de duplicados usando un dict _last_event_ts separado por huella digital que almacena la última marca de tiempo del evento vista para esa huella digital.
# Para evitar más confusión, pasar a una implementación más sencilla: mantener un dict auxiliar _last_event_ts por fp.
# Debido a que lo anterior se volvió desordenado dentro de ingest, reescribiremos el método ingest limpiamente usando estado auxiliar.
Reescribiendo el módulo con la estructura corregida a continuación (reemplazo de archivo único)
Empezamos la nueva implementación aquí
import heapq
class AlertEngine:
"""Implementación de AlertEngine (reescritura limpia).
Ver la docstring de la clase anterior para notas de comportamiento.
"""
def __init__(self, config: dict, now: float):
try:
self.window_seconds = float(config["window_seconds"])
self.max_alerts_per_window = int(config["max_alerts_per_window"])
self.dedup_seconds = float(config["dedup_seconds"])
self.severity_floor = str(config["severity_floor"]).lower()
if self.severity_floor not in _SEVERITY_ORDER:
raise KeyError
except Exception:
raise ValueError("configuración inválida")
be = config.get("burst_escalation")
if be is not None:
try:
self.burst_count = int(be["count"])
self.burst_within = float(be["within_seconds"])
if self.burst_count <= 0 or self.burst_within <= 0:
raise ValueError
except Exception:
raise ValueError("configuración de burst_escalation inválida")
else:
self.burst_count = None
self.burst_within = None
self._now_watermark = float(now)
self._fp_store: Dict[str, FPEntry] = {}
self._last_event_ts: Dict[str, float] = {}
self._key_emissions: Dict[Tuple[str, str], deque] = defaultdict(deque)
self._total_ingested = 0
self._total_emitted = 0
self._total_suppressed = 0
def _clamp_now(self, now: float) -> float:
if now < self._now_watermark:
return self._now_watermark
self._now_watermark = now
return now
def _prune_key_emissions(self, key: Tuple[str, str], now: float) -> None:
q = self._key_emissions.get(key)
if not q:
return
cutoff = now - self.window_seconds
while q and q[0] < cutoff:
q.popleft()
def _emit_alert(self, kind: str, fp: str, entry: FPEntry, count: int) -> dict:
alert = {
"kind": kind,
"fingerprint": fp,
"key": (entry.service, entry.severity),
"first_seen": entry.first_seen,
"last_seen": entry.last_seen,
"count": count,
"sample_message": entry.sample_message,
}
self._total_emitted += 1
return alert
def ingest(self, event: dict, now: float) -> List[dict]:
self._total_ingested += 1
now = self._clamp_now(float(now))
alerts: List[dict] = []
# validar
if not isinstance(event, dict):
return []
for k in ("timestamp", "severity", "service", "message"):
if k not in event:
return []
try:
ev_ts = float(event["timestamp"])
severity = str(event["severity"]).lower()
service = str(event["service"])
message = str(event["message"])
except Exception:
return []
if severity not in _SEVERITY_ORDER:
return []
if _SEVERITY_ORDER[severity] < _SEVERITY_ORDER[self.severity_floor]:
return []
fp, normalized = _fingerprint_for(service, severity, message)
sample_message = message[:MAX_SAMPLE_LEN]
prev_ts = self._last_event_ts.get(fp)
is_duplicate = prev_ts is not None and (ev_ts - prev_ts) <= self.dedup_seconds
# crear o actualizar entrada
entry = self._fp_store.get(fp)
if entry is None:
entry = FPEntry(
service=service,
severity=severity,
normalized=normalized,
first_seen=ev_ts,
last_seen=ev_ts,
count=1,
sample_message=sample_message,
suppressed_since_emit=0,
burst_times=deque(),
last_emitted=None,
)
self._fp_store[fp] = entry
else:
entry.count += 1
# mantener first_seen como está
entry.last_seen = ev_ts
# mantener una muestra estable, no sobrescribir, pero si está vacía, establecer
if not entry.sample_message:
entry.sample_message = sample_message
# actualizar last_event_ts para verificaciones de deduplicación la próxima vez
self._last_event_ts[fp] = ev_ts
# Evacuar si es necesario
if len(self._fp_store) > MAX_FINGERPRINTS:
# desalojo determinista por last_seen, first_seen
items = sorted(self._fp_store.items(), key=lambda kv: (kv[1].last_seen, kv[1].first_seen))
to_evict = len(self._fp_store) - MAX_FINGERPRINTS
for i in range(to_evict):
del self._fp_store[items[i][0]]
if is_duplicate:
# Suprimir duplicado
entry.suppressed_since_emit += 1
self._total_suppressed += 1
# rastrear tiempos de ráfaga usando el tiempo de procesamiento 'now'
if self.burst_count is not None:
bt = entry.burst_times
bt.append(now)
cutoff = now - self.burst_within
while bt and bt[0] < cutoff:
bt.popleft()
if len(bt) > self.burst_count:
# emitir escalada: omitir límite de velocidad, restablecer contador de ráfaga
alert = self._emit_alert("escalation", fp, entry, entry.suppressed_since_emit)
alerts.append(alert)
entry.suppressed_since_emit = 0
bt.clear()
entry.last_emitted = now
return alerts
# No es un duplicado: candidato para nueva alerta. Comprobar límite de velocidad por clave.
key = (service, severity)
self._prune_key_emissions(key, now)
q = self._key_emissions[key]
if len(q) < self.max_alerts_per_window:
# permitido
q.append(now)
alert = self._emit_alert("new", fp, entry, 1)
alerts.append(alert)
entry.last_emitted = now
entry.suppressed_since_emit = 0
else:
# límite de velocidad: suprimir y contar
entry.suppressed_since_emit += 1
self._total_suppressed += 1
# no emitir inmediatamente; aparecerá en flush
return alerts
def flush(self, now: float) -> List[dict]:
now = self._clamp_now(float(now))
alerts: List[dict] = []
# Para cada huella digital, si la ventana de deduplicación se cerró y hay eventos suprimidos, emitir rate_limit_notice
to_delete = []
for fp, entry in list(self._fp_store.items()):
# si last_seen es anterior a dedup_seconds
if now - entry.last_seen >= self.dedup_seconds and entry.suppressed_since_emit > 0:
alert = self._emit_alert("rate_limit_notice", fp, entry, entry.suppressed_since_emit)
alerts.append(alert)
entry.suppressed_since_emit = 0
entry.last_emitted = now
return alerts
def stats(self) -> dict:
return {
"total_ingested": self._total_ingested,
"total_emitted": self._total_emitted,
"total_suppressed": self._total_suppressed,
"active_fingerprints": len(self._fp_store),
}
test_alert_engine.py
import unittest
class TestAlertEngine(unittest.TestCase):
def setUp(self):
self.config = {
"window_seconds": 60.0,
"max_alerts_per_window": 2,
"dedup_seconds": 10.0,
"severity_floor": "info",
"burst_escalation": {"count": 3, "within_seconds": 5.0},
}
self.engine = AlertEngine(self.config, now=0.0)
def test_dedup_collapsing(self):
now = 1.0
e1 = {"timestamp": now, "severity": "info", "service": "svc", "message": "user 123 logged in"}
out = self.engine.ingest(e1, now)
self.assertEqual(len(out), 1)
fp = out[0]["fingerprint"]
# duplicado dentro de dedup_seconds
e2 = {"timestamp": now + 2, "severity": "info", "service": "svc", "message": "user 456 logged in"}
out2 = self.engine.ingest(e2, now + 2)
self.assertEqual(out2, [])
stats = self.engine.stats()
self.assertEqual(stats["total_suppressed"], 1)
def test_rate_limiting_at_boundary(self):
now = 10.0
# permitir dos alertas por ventana
for i in range(2):
e = {"timestamp": now + i, "severity": "error", "service": "s", "message": f"msg{i}"}
out = self.engine.ingest(e, now + i)
self.assertEqual(len(out), 1)
# el tercero debería ser suprimido
e3 = {"timestamp": now + 3, "severity": "error", "service": "s", "message": "msg3"}
out3 = self.engine.ingest(e3, now + 3)
self.assertEqual(out3, [])
stats = self.engine.stats()
self.assertEqual(stats["total_suppressed"], 1)
# después de que pase la ventana, flush debería producir un resumen de los suprimidos
out_flush = self.engine.flush(now + 70)
self.assertTrue(any(a["kind"] == "rate_limit_notice" for a in out_flush))
def test_burst_escalation(self):
now = 100.0
e = {"timestamp": now, "severity": "warn", "service": "svcB", "message": "hit 1"}
out = self.engine.ingest(e, now)
self.assertEqual(len(out), 1)
fp = out[0]["fingerprint"]
# producir duplicados suprimidos rápidamente para exceder el recuento de ráfagas (count=3)
for i in range(4):
ed = {"timestamp": now + 1 + i, "severity": "warn", "service": "svcB", "message": f"hit {10+i}"}
res = self.engine.ingest(ed, now + 1 + i)
if res:
# uno de ellos debería ser una escalada cuando se supera el umbral
kinds = {r["kind"] for r in res}
self.assertIn("escalation", kinds)
break
else:
self.fail("escalation not emitted")
def test_flush_idempotency(self):
now = 200.0
e = {"timestamp": now, "severity": "error", "service": "sF", "message": "a"}
self.engine.ingest(e, now)
# suprimir el siguiente por límite de velocidad
# crear dos alertas más para alcanzar el límite
self.engine.ingest({"timestamp": now+1, "severity": "error", "service": "sF", "message": "b"}, now+1)
self.engine.ingest({"timestamp": now+2, "severity": "error", "service": "sF", "message": "c"}, now+2)
out1 = self.engine.flush(now+30)
out2 = self.engine.flush(now+31)
# el segundo flush debería ser idempotente
self.assertEqual(out1, out2)
def test_out_of_order_timestamps(self):
now = 300.0
e1 = {"timestamp": now, "severity": "info", "service": "oo", "message": "x1"}
out1 = self.engine.ingest(e1, now)
self.assertEqual(len(out1), 1)
# suministrar un 'now' anterior (evento tardío). El motor ajusta el tiempo y no debería retroceder.
out2 = self.engine.ingest({"timestamp": now-50, "severity": "info", "service": "oo", "message": "x2"}, now-50)
# no debe fallar y no emitir fuera de orden
self.assertIsInstance(out2, list)
def test_severity_floor_filtering(self):
now = 400.0
e = {"timestamp": now, "severity": "debug", "service": "sD", "message": "dmsg"}
out = self.engine.ingest(e, now)
self.assertEqual(out, [])
def test_malformed_input(self):
now = 500.0
out = self.engine.ingest({"severity": "info"}, now)
self.assertEqual(out, [])
out2 = self.engine.ingest("not a dict", now)
self.assertEqual(out2, [])
def test_eviction_behavior(self):
# Reducir MAX_FINGERPRINTS para la prueba parcheando temporalmente la variable global
global MAX_FINGERPRINTS
old = MAX_FINGERPRINTS
MAX_FINGERPRINTS = 5
try:
eng = AlertEngine(self.config, now=0.0)
for i in range(10):
e = {"timestamp": i, "severity": "info", "service": f"svc{i}", "message": "m"}
eng.ingest(e, float(i))
stats = eng.stats()
self.assertLessEqual(stats["active_fingerprints"], 5)
finally:
MAX_FINGERPRINTS = old
if name == 'main':
unittest.main()
NOTA DE DISEÑO
"""
Estructuras de datos: un dict con clave por huella digital almacena dataclasses FPEntry (primera/última vista, recuentos, mensaje de muestra, contador suprimido, marcas de tiempo de ráfaga). Las colas de emisión por clave (servicio, severidad) mantienen marcas de tiempo de emisión recientes en deques para poda de ventana deslizante O(1). Un mapa de marca de tiempo del último evento se utiliza para detectar duplicados de forma determinista.
Complejidad temporal: ingest es O(1) en promedio — el hashing y el trabajo de expresiones regulares sobre el mensaje (lineal en la longitud del mensaje), las operaciones de deque para emisión y ráfaga son O(1) amortizado. Flush es O(N) en huellas dactilares activas.
La evacuación de memoria es determinista: cuando el almacén de huellas dactilares excede MAX_FINGERPRINTS (es un límite global), el motor evacua las huellas dactilares más antiguas por (last_seen, first_seen). La política de desalojo: LRU determinista por last_seen con desempate por first_seen; esto limita la memoria y es predecible. Contrapartida: la evacuación utiliza ordenación cuando se excede el límite, lo que es O(M log M) para M elementos; MAX_FINGERPRINTS mantiene esto limitado. Una contrapartida aceptada: la detección de duplicados utiliza la marca de tiempo del evento y una política de marca de agua para 'now' no monótono (ajustando 'now' al máximo visto). Esto simplifica las garantías de ordenación (no se emiten alertas fuera del orden interno) pero significa que los tiempos suministrados muy tardíos se tratan como si hubieran ocurrido en la marca de agua, alterando ligeramente la fidelidad temporal.
"""