alert_engine.py
"""
AlertEngine: ストリーミングログの重複排除とレート制限付きアラートエンジン。
標準ライブラリのみを使用し、バックグラウンドスレッドは使用しません。外部から供給される「現在時刻」を使用します。
"""
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
重み付けマッピング
_SEVERITY_ORDER = {
"debug": 10,
"info": 20,
"warn": 30,
"warning": 30, # 別名を受け入れる
"error": 40,
"critical": 50,
}
MAX_SAMPLE_LEN = 200
MAX_FINGERPRINTS = 100000 # 決定論的なメモリキャップ。設計ノートを参照
_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:
# UUID、IP、長い16進数ID、数字の連続を置換します。順序が重要です。
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 # バースト検出のための抑制された重複のタイムスタンプ(float)
last_emitted: Optional[float]
class AlertEngine:
"""AlertEngine(config: dict, now: float)
configキー:
- window_seconds: float
- max_alerts_per_window: int
- dedup_seconds: float
- severity_floor: debug, info, warn, error, critical のいずれか
- burst_escalation: オプションの辞書 {"count": int, "within_seconds": float}
非単調なnowに対するポリシー: エンジンは内部ウォーターマーク(検出された最大now)を維持します。
ingest()またはflush()がウォーターマークより前のnowで呼び出された場合、
エンジンは提供されたnowをウォーターマークとして扱います(つまり、時間を前方にクランプします)。
これにより、発行されるアラートが時間順に逆転することなく、
遅延イベントに対しても決定論的な動作が保証されます。
追い出しポリシー: メモリを制限するために、決定論的なキャップMAX_FINGERPRINTSが適用されます。
超過した場合、エンジンは最も古い(最小の)last_seenを持つフィンガープリントを追い出します。
同点の場合はfirst_seenで解決します。これは決定論的であり、無制限の成長を防ぎます。
"""
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("無効な設定: 必須キー 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("無効な burst_escalation 設定")
else:
self.burst_count = None
self.burst_within = None
# 内部状態
self._now_watermark = float(now)
self._fp_store: Dict[str, FPEntry] = {}
# キーごと(service, severity)の発生タイムスタンプキュー(スライディングウィンドウ用)
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 _ensure_eviction(self) -> None:
if len(self._fp_store) <= MAX_FINGERPRINTS:
return
# 決定論的な追い出し: 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):
k = items[i][0]
del self._fp_store[k]
def ingest(self, event: dict, now: float) -> List[dict]:
"""単一イベントを処理し、発行されたアラートのリストを返します(空の場合があります)。
不正な形式の入力に対するポリシー: エンジンはイベントを取り込んだものとしてカウントしますが、
必要なフィールドが欠落しているか、重大度が無効な場合はドロップして空のリストを返します。
"""
self._total_ingested += 1
now = self._clamp_now(float(now))
alerts: List[dict] = []
# 最小限の検証
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:
# 明示的な処理: 重大度が不明な場合はイベントをドロップ
return []
# 重み付けフロア
if _SEVERITY_ORDER[severity] < _SEVERITY_ORDER[self.severity_floor]:
# 他のすべての処理の前にドロップされる
return []
# 正規化とフィンガープリント化
fp, normalized = _fingerprint_for(service, severity, message)
# サンプルメッセージをフィンガープリントに影響を与えずに切り捨てる
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:
# タイムスタンプ/カウンターの更新
entry.count += 1
entry.last_seen = ev_ts
# 最初のサンプルメッセージを保持
if len(entry.sample_message) < MAX_SAMPLE_LEN:
# 初期メッセージが短かった場合にサンプルを埋めることを試みる
entry.sample_message = (entry.sample_message + " | " + sample_message)[:MAX_SAMPLE_LEN]
# 重複チェック: last_seen(イベントタイムスタンプを使用)からdedup_seconds以内であれば重複とみなす
is_duplicate = (ev_ts - entry.last_seen) <= self.dedup_seconds if entry.count > 1 else False
# 注意: last_seenはすでにev_tsに設定されています。重複ロジックでは前のlast_seenを使用する必要があります。
# 正しく実装するには再計算します。count==1の場合 -> 重複ではない。それ以外の場合、ev_ts - prev_last_seen <= dedup_seconds
if entry.count == 1:
is_duplicate = False
else:
# 前のlast_seenが必要です。ev_ts - (entry.last_seen または ev_ts) をチェックして近似しますが、last_seenは更新済みです。
# より簡単な方法: 現在のイベントの前にlast_seenを一時的に保存しておけば、それを参照できます。
# 正しく修正するには、上記の順序を修正する必要があります。
pass
# 上記のロジックは、last_seenを早期に更新したため、扱いにくいです。修正するには、適切な順序で再実行します。
# 再構築: 前のエントリのスナップショットを取得します。
# ----- 現在時刻の処理におけるイベントコアの再実装(正しい時間管理を使用)-----
# このエントリに対する以前の更新の影響を元に戻す
# 再構築または調整。保存されたフィールドから再構築します。
stored_entry = self._fp_store[fp]
# stored_entry.countがこのイベントを含むと仮定して、以前のカウントを取得するには1を引きます。
previous_count = stored_entry.count - 1
# previous_last_seenは現在信頼できません。しかし、重複検出は現在時刻(処理時刻)を使用して行うことができます。
# ポリシーを使用: 重複は、now - stored_entry.last_seen <= dedup_seconds の場合に検出されます。
# ここでstored_entry.last_seenは現在のイベントの前に更新されたものです。
# これを達成するために、フィンガープリントごとに最後のイベントタイムスタンプを格納する補助的な辞書_last_event_tsを維持します。
# ingestメソッド内の上記が混乱したため、ヘルパー状態を使用してクリーンに書き直します。
ingestメソッドの混乱を避けるため、モジュール全体をクリーンに書き直します。
ここから新しい実装を開始します
import heapq
class AlertEngine:
"""AlertEngineの実装(クリーンな書き直し)。
動作に関する注記については、以前のクラスのdocstringを参照してください。
"""
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("無効な設定")
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("無効な burst_escalation 設定")
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] = []
# 検証
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
# エントリの作成または更新
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
# first_seenはそのまま保持
entry.last_seen = ev_ts
# 安定したサンプルを保持し、上書きしない。ただし、空の場合は設定する。
if not entry.sample_message:
entry.sample_message = sample_message
# 次回の重複チェックのためにlast_event_tsを更新
self._last_event_ts[fp] = ev_ts
# 必要に応じて追い出し
if len(self._fp_store) > MAX_FINGERPRINTS:
# 決定論的な追い出し: 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:
# 重複を抑制
entry.suppressed_since_emit += 1
self._total_suppressed += 1
# バーストタイムを処理時刻 '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:
# エスカレーションを発行: レート制限をバイパスし、バーストカウンターをリセット
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
# 重複ではない: 新しいアラートの候補。キーごとのレート制限を確認。
key = (service, severity)
self._prune_key_emissions(key, now)
q = self._key_emissions[key]
if len(q) < self.max_alerts_per_window:
# 許可
q.append(now)
alert = self._emit_alert("new", fp, entry, 1)
alerts.append(alert)
entry.last_emitted = now
entry.suppressed_since_emit = 0
else:
# レート制限: 抑制し、カウントする
entry.suppressed_since_emit += 1
self._total_suppressed += 1
# すぐには発行しない。flush時に表示される。
return alerts
def flush(self, now: float) -> List[dict]:
now = self._clamp_now(float(now))
alerts: List[dict] = []
# 各フィンガープリントについて、重複ウィンドウが閉じ、抑制されたイベントがある場合、rate_limit_noticeを発行する
to_delete = []
for fp, entry in list(self._fp_store.items()):
# last_seenが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"]
# 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
# ウィンドウごとに2つのアラートを許可
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)
# 3番目は抑制されるべき
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)
# ウィンドウ経過後、flushは抑制されたアラートのサマリーを生成するはず
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"]
# バーストカウント(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:
# そのうちの1つは、しきい値を超えたときの増幅であるべき
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)
# レート制限により次のイベントを抑制
# 2つの追加アラートを作成して制限に達する
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)
# 2回目のflushは冪等であるべき
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)
# より早いnowを供給(遅延イベント)。エンジンは時間をクランプし、後退しないはずです。
out2 = self.engine.ingest({"timestamp": now-50, "severity": "info", "service": "oo", "message": "x2"}, now-50)
# クラッシュせず、順序外で発行しないこと
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):
# テストのためにMAX_FINGERPRINTSを減らすためにグローバルを一時的にパッチする
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()
DESIGN NOTE
"""
データ構造: フィンガープリントをキーとする辞書は、FPEntry dataclass(初回/最終表示、カウント、
サンプルメッセージ、抑制カウンター、バーストタイムスタンプ)を格納します。キーごと(service,severity)の
発生キューは、O(1)スライディングウィンドウのプルーニングのために、最近の発生タイムスタンプをdequeに保持します。
最後のイベントタイムスタンプマップは、決定論的に重複を検出するために使用されます。
時間計算量: ingestは平均O(1)です。ハッシュと正規表現はメッセージ(メッセージ長に線形)で機能し、
発生とバーストのためのdeque操作は償却O(1)です。flushはアクティブなフィンガープリントのO(N)です。
メモリ追い出しは決定論的です。フィンガープリントストアがMAX_FINGERPRINTS(グローバルキャップ)を超えた場合、
エンジンは最も古いフィンガープリントを(last_seen, first_seen)で追い出します。
追い出しポリシー: last_seenによる決定論的なLRUライクな動作で、first_seenで同点を解決します。これによりメモリが制限され、予測可能です。
トレードオフ: キャップを超えた場合の追い出しではソートを使用し、これはM個のアイテムに対してO(M log M)です。MAX_FINGERPRINTSはこの値を制限内に保ちます。
受け入れられたトレードオフ: 重複検出はイベントタイムスタンプとウォーターマークポリシーを使用して非単調な「now」(最大seenへのクランプ)を処理します。
これにより順序保証が簡素化されます(内部順序外でアラートが発行されない)。しかし、非常に遅延した供給時間は
ウォーターマーク発生時に発生したかのように扱われ、時間的忠実度がわずかに変化します。
"""