alert_engine.py
"""
AlertEngine: motor de deduplicação de logs em streaming e de alerta com taxa limitada.
Biblioteca padrão pura, sem threads em segundo plano. Usa tempos 'now' fornecidos 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
Mapeamento de severidade
_SEVERITY_ORDER = {
"debug": 10,
"info": 20,
"warn": 30,
"warning": 30, # aceita alternativa
"error": 40,
"critical": 50,
}
MAX_SAMPLE_LEN = 200
MAX_FINGERPRINTS = 100000 # limite de memória determinístico; veja nota de design
_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:
# Substitui UUIDs, IPs, IDs hexadecimais longos, depois sequências de dígitos. A ordem 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 # timestamps (floats) de duplicatas suprimidas para detecção de rajadas
last_emitted: Optional[float]
class AlertEngine:
"""AlertEngine(config: dict, now: float)
chaves de config:
- window_seconds: float
- max_alerts_per_window: int
- dedup_seconds: float
- severity_floor: um de debug, info, warn, error, critical
- burst_escalation: dict opcional {"count": int, "within_seconds": float}
Política para 'now' não monotônico: o motor mantém um marcador interno (o 'now' máximo visto).
Se ingest() ou flush() for chamado com um 'now' anterior ao marcador, o motor trata o 'now' fornecido
como o marcador (ou seja, ele limita o tempo para frente).
Isso garante que os alertas emitidos nunca sejam produzidos fora de ordem por tempo e torna
o comportamento determinístico para eventos tardios.
Política de expurgo: para limitar a memória, um limite determinístico MAX_FINGERPRINTS é imposto.
Quando excedido, o motor expulsa fingerprints com o mais antigo (menor) last_seen,
quebrando empates por first_seen. Isso é determinístico e evita crescimento ilimitado.
"""
def __init__(self, config: dict, now: float):
# validação mínima da 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("config inválida: chaves necessárias 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("config 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] = {}
# fila por chave (service,severity) de timestamps de emissão para janela deslizante
self._key_emissions: Dict[Tuple[str, str], deque] = defaultdict(deque)
# estatí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 atrasado; limitar para frente ao marcador
return self._now_watermark
self._now_watermark = now
return now
def _ensure_eviction(self) -> None:
if len(self._fp_store) <= MAX_FINGERPRINTS:
return
# expurgo determinístico: ordenar por last_seen, depois 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]:
"""Processa um único evento e retorna uma lista de alertas emitidos (pode estar vazia).
Política para entrada malformada: o motor conta o evento como ingerido, mas o descarta
e retorna uma lista vazia quando os campos necessários estão ausentes ou a severidade é inválida.
"""
self._total_ingested += 1
now = self._clamp_now(float(now))
alerts: List[dict] = []
# validação 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:
# tratamento explícito: severidade desconhecida -> descarta evento
return []
# limite de severidade
if _SEVERITY_ORDER[severity] < _SEVERITY_ORDER[self.severity_floor]:
# descartado antes de qualquer outro processamento
return []
# normaliza e gera fingerprint
fp, normalized = _fingerprint_for(service, severity, message)
# garante que sample_message seja truncado sem afetar o fingerprint
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:
# atualiza timestamps/contadores
entry.count += 1
entry.last_seen = ev_ts
# mantém a primeira sample_message
if len(entry.sample_message) < MAX_SAMPLE_LEN:
# tenta preencher a amostra se a inicial foi curta
entry.sample_message = (entry.sample_message + " | " + sample_message)[:MAX_SAMPLE_LEN]
# verificação de dedup: se dentro de dedup_seconds de last_seen (usando timestamp do evento), trata como duplicata
is_duplicate = (ev_ts - entry.last_seen) <= self.dedup_seconds if entry.count > 1 else False
# Nota: last_seen já foi definido para ev_ts; para lógica de duplicata, devemos usar o last_seen anterior
# Para implementar corretamente, recalcular: se count==1 -> não é duplicata; senão se ev_ts - prev_last_seen <= dedup_seconds
if entry.count == 1:
is_duplicate = False
else:
# precisamos do last_seen anterior: aproximar verificando ev_ts - (entry.last_seen ou ev_ts) mas atualizamos last_seen.
# Mais simples: consideramos duplicatas como eventos cujo event.timestamp está dentro de dedup_seconds do timestamp do evento anterior.
# Como atualizamos, não podemos recuperar o anterior; corrigir mantendo last_seen_old temporariamente antes.
pass
# A lógica acima é estranha porque atualizamos last_seen precocemente. Para corrigir, refazer com rastreamento anterior adequado.
# Recalcular com armazenamento explícito.
# Refazer: obter snapshot da entrada anterior
prev_last_seen = entry.last_seen
prev_count = entry.count
# Recalcular corretamente: se prev_count > 1, o last_seen anterior não é acessível; para evitar complexidade,
# trataremos a detecção de duplicatas usando agora (tempo de processamento) em relação a entry.last_seen armazenado antes do evento atual.
# Para conseguir isso, salvamos last_seen igual a ev_ts acima; mas se previous_count>=1, o last_seen anterior é igual a entry.first_seen para count==1,
# ou algo mais. Essa complexidade surge porque atualizamos campos prematuramente.
# Para simplificar e tornar determinístico: consideraremos duplicatas se now - entry.last_seen <= dedup_seconds usando last_seen antes de atualizar.
# Como não podemos recuperar o last_seen anterior, manteremos a detecção de dedup usando um _last_event_ts auxiliar por fingerprint, armazenando o último timestamp do evento visto para esse fingerprint.
# Para evitar mais confusão, mover para uma implementação mais direta: manter um dicionário auxiliar _last_event_ts por fp.
# Como o acima ficou confuso dentro do ingest, reescreveremos o método ingest de forma limpa usando estado auxiliar.
Reescrevendo o módulo com estrutura corrigida abaixo (substituição de arquivo único)
Começa a nova implementação aqui
import heapq
class AlertEngine:
"""Implementação do AlertEngine (reescrita limpa).
Veja a docstring da classe anterior para notas de comportamento.
"""
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("config 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("config 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] = []
# valida
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
# cria ou atualiza 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
# mantém first_seen como está
entry.last_seen = ev_ts
# mantém uma amostra estável, não sobrescreve, mas se estiver vazia, define
if not entry.sample_message:
entry.sample_message = sample_message
# atualiza _last_event_ts para verificações de dedup na próxima vez
self._last_event_ts[fp] = ev_ts
# Expulsa se necessário
if len(self._fp_store) > MAX_FINGERPRINTS:
# expurgo determinístico 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:
# Suprime duplicata
entry.suppressed_since_emit += 1
self._total_suppressed += 1
# rastreia tempos de rajada usando o tempo de processamento '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:
# emite escalonamento: ignora limite de taxa, reseta contador de rajada
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
# Não é uma duplicata: candidato a novo alerta. Verifica limite de taxa por chave.
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:
# limitado pela taxa: suprime e conta
entry.suppressed_since_emit += 1
self._total_suppressed += 1
# não emite imediatamente; aparecerá no flush
return alerts
def flush(self, now: float) -> List[dict]:
now = self._clamp_now(float(now))
alerts: List[dict] = []
# Para cada fingerprint, se a janela de dedup fechou e há eventos suprimidos, emite rate_limit_notice
to_delete = []
for fp, entry in list(self._fp_store.items()):
# se last_seen mais antigo que 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"]
# duplicata 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
# permite dois alertas por janela
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)
# o terceiro deve 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)
# após a janela passar, o flush deve produzir um resumo para os 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"]
# produz duplicatas suprimidas rapidamente para exceder a contagem de rajada (contagem=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:
# um deles deve ser escalonamento quando o limite for atingido
kinds = {r["kind"] for r in res}
self.assertIn("escalation", kinds)
break
else:
self.fail("escalonamento não emitido")
def test_flush_idempotency(self):
now = 200.0
e = {"timestamp": now, "severity": "error", "service": "sF", "message": "a"}
self.engine.ingest(e, now)
# suprime o próximo por limitação de taxa
# cria mais dois alertas para atingir o limite
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)
# o segundo flush deve 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)
# fornece 'now' anterior (evento atrasado). O motor limita o tempo e não deve voltar.
out2 = self.engine.ingest({"timestamp": now-50, "severity": "info", "service": "oo", "message": "x2"}, now-50)
# não deve falhar e não emitir fora de ordem
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):
# Reduz MAX_FINGERPRINTS para teste, modificando temporariamente o 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 DESIGN
"""
Estruturas de dados: um dicionário com chave por fingerprint armazena dataclasses FPEntry (primeira/última vista, contagens, mensagem de amostra, contador suprimido, timestamps de rajada). Filas de emissão por chave (service,severity) mantêm timestamps de emissão recentes em deques para poda de janela deslizante O(1). Um mapa de timestamp do último evento é usado para detectar duplicatas deterministicamente.
Complexidade de tempo: ingest é O(1) em média — hashing e trabalho de regex na mensagem (linear no comprimento da mensagem), operações de deque para emissão e rajada são O(1) amortizado. flush é O(N) em fingerprints ativos.
Política de expurgo: determinística semelhante a LRU por last_seen com desempate por first_seen; isso limita a memória e é previsível. Contrapartida: o expurgo usa ordenação quando o limite é excedido, o que é O(M log M) para M itens; MAX_FINGERPRINTS mantém isso limitado. Uma contrapartida aceita: a detecção de dedup usa o timestamp do evento e uma política de watermark para 'now' não monotônico (limitando 'now' ao máximo visto). Isso simplifica as garantias de ordenação (nenhum alerta emitido fora da ordem interna), mas significa que tempos fornecidos muito atrasados são tratados como se tivessem ocorrido no watermark, alterando ligeiramente a fidelidade temporal.
"""